feat: add interview-prep topic group with 250 questions (6 subtopics, 5 question types)
Deploy Examination / deploy (push) Successful in 10s

Subtopics:
- distributed-microservice: 45 questions (分布式微服务架构)
- message-queue: 45 questions (消息队列)
- k8s-observability: 45 questions (K8s与可观测性)
- go-java-concurrency: 45 questions (Go/Java并发模型)
- database-advanced: 35 questions (数据库进阶)
- ai-engineering: 35 questions (AI工程实践)

Question types: single_choice, true_false, fill_blank, short_answer, code_reading
This commit is contained in:
2026-09-09 16:36:27 +08:00
parent 515acbcc7f
commit 0f68a64829
38 changed files with 6835 additions and 2 deletions
@@ -0,0 +1,303 @@
{
"topic": "message-queue",
"type": "code_reading",
"schema_version": "1.0.0",
"generated": "2026-09-09T16:00:00+08:00",
"questions": [
{
"id": "cr-001",
"type": "code_reading",
"difficulty": 4,
"tags": [
"message-queue",
"kafka",
"consumer-group",
"partition"
],
"question": "阅读以下 Go 语言实现的 Kafka 消费者 Rebalance 监听器代码,分析其工作流程。",
"language": "go",
"code": "package kafka\n\nimport (\n \"log\"\n \"sync\"\n)\n\n// ConsumerRebalanceHandler 实现了 Kafka 消费者的 Rebalance 回调\n\ntype ConsumerRebalanceHandler struct {\n consumerGroup string\n partitions map[string][]int32 // topic -> partitions\n mu sync.RWMutex\n rebalanceCount int\n}\n\n// OnPartitionsAssigned 当分配到新分区时触发\nfunc (h *ConsumerRebalanceHandler) OnPartitionsAssigned(\n consumerGroup string,\n topicPartitions map[string][]int32,\n) {\n h.mu.Lock()\n defer h.mu.Unlock()\n \n // 1. 记录新分配的分区\n for topic, partitions := range topicPartitions {\n h.partitions[topic] = partitions\n }\n \n // 2. 重新初始化分区级别的消费位点缓存\n for topic, partitions := range topicPartitions {\n for _, p := range partitions {\n log.Printf(\"[ASSIGNED] group=%s topic=%s partition=%d\",\n consumerGroup, topic, p)\n }\n }\n \n h.rebalanceCount++\n log.Printf(\"Rebalance #%d 完成,共分配 %d 个分区\",\n h.rebalanceCount, countPartitions(topicPartitions))\n}\n\n// OnPartitionsRevoked 当分区被回收时触发\nfunc (h *ConsumerRebalanceHandler) OnPartitionsRevoked(\n consumerGroup string,\n topicPartitions map[string][]int32,\n) {\n h.mu.Lock()\n defer h.mu.Unlock()\n \n // 1. 先提交当前消费位点(确保不丢消息)\n for topic, partitions := range topicPartitions {\n for _, p := range partitions {\n log.Printf(\"[REVOKED] 提交位点: group=%s topic=%s partition=%d\",\n consumerGroup, topic, p)\n }\n }\n \n // 2. 清理分区状态\n for topic := range topicPartitions {\n delete(h.partitions, topic)\n }\n}\n\n// OnPartitionsLost 当分区不可用时触发(比 Revoked 更激进)\nfunc (h *ConsumerRebalanceHandler) OnPartitionsLost(\n consumerGroup string,\n topicPartitions map[string][]int32,\n) {\n // 不提交位点,直接清理\n h.mu.Lock()\n defer h.mu.Unlock()\n \n for topic := range topicPartitions {\n delete(h.partitions, topic)\n }\n log.Printf(\"[LOST] 分区丢失,未提交位点\")\n}",
"sub_questions": [
{
"index": 1,
"type": "single_choice",
"question": "在 OnPartitionsRevoked 回调中,为什么要先提交消费位点再清理分区状态?",
"options": {
"A": "为了提高吞吐量,将位点提交和状态清理合并执行",
"B": "确保在分区被其他消费者接管前,当前已消费的进度不会丢失",
"C": "Kafka 协议要求在 revoke 时必须提交位点,否则会报错",
"D": "为了清理 Broker 端的位点缓存,释放存储空间"
},
"answer": "B",
"explanation": "Rebalance 的本质是将分区重新分配给其他消费者。如果不先提交位点就清理状态,新的消费者将从上次提交的位点开始消费,导致已处理的消息被重复消费。先提交位点可以保证新消费者从正确的位置继续,避免消息丢失或大量重复。Kafka 协议本身并不强制要求 revoke 时提交位点(选项 C 错误),这是一个最佳实践。"
},
{
"index": 2,
"type": "single_choice",
"question": "OnPartitionsLost 与 OnPartitionsRevoked 的核心区别是什么?",
"options": {
"A": "OnPartitionsLost 会触发更重的 Rebalance 流程",
"B": "OnPartitionsLost 不提交位点,适用于分区分配异常丢失的场景;OnPartitionsRevoked 是正常的 Rebalance 流程",
"C": "两者功能完全相同,只是调用时机不同",
"D": "OnPartitionsLost 是客户端行为,OnPartitionsRevoked 是服务端行为"
},
"answer": "B",
"explanation": "在 Kafka 的 Cooperative Rebalance 协议中,OnPartitionsLost 表示分区被异常剥夺(如 Consumer 挂了但 session 超时未过),此时无法安全提交位点(可能已在别处提交),所以直接清理状态。OnPartitionsRevoked 是正常的 Rebalance 流程,消费者有机会提交位点后再释放分区。代码中 OnPartitionsLost 也注释了「未提交位点」,印证了这一设计意图。"
},
{
"index": 3,
"type": "short_answer",
"question": "代码中使用 sync.RWMutex 而非 sync.Mutex 的原因是什么?请结合消费场景分析其性能优势。",
"answer": "OnPartitionsAssigned 和 OnPartitionsRevoked 只在 Rebalance 发生时被调用(低频),而消费者主线程会频繁读取 h.partitions 来判断某分区是否仍由自己负责(高频)。RWMutex 允许多个读操作并发执行,只有写操作(Rebalance 回调)才需要独占锁,因此用 RWMutex 可以避免高频读操作之间的锁竞争,显著提升正常消费路径的吞吐量。",
"keywords": [
"读写锁",
"高频读",
"低频写",
"Rebalance",
"消费主线程",
"并发"
],
"scoring_rubric": "答出「Rebalance 回调低频写、消费路径高频读」的核心场景差异得 3 分;提到 RWMutex 允许并发读得 2 分;提到避免锁竞争或性能提升得 1 分。满分 6 分,4 分及以上通过。"
}
],
"explanation": "本题考查 Kafka 消费者 Rebalance 机制的核心代码逻辑。Rebalance 是 Consumer Group 模型的关键机制,当消费者加入或离开组时,分区会被重新分配。理解三个回调(Assigned/Revoked/Lost)的调用时机和职责差异,是正确实现消费者的基础。重点掌握:(1) Revoked 时先提交位点防止消息重复;(2) Lost 与 Revoked 的语义差异;(3) 并发控制在消费路径中的应用。",
"source": null,
"related": []
},
{
"id": "cr-002",
"type": "code_reading",
"difficulty": 5,
"tags": [
"message-queue",
"rocketmq",
"transactional-message",
"durable"
],
"question": "阅读以下 Java 代码,分析 RocketMQ 事务消息的发送与回查机制实现。",
"language": "java",
"code": "package com.example.mq.transaction;\n\nimport org.apache.rocketmq.client.producer.*;\nimport org.apache.rocketmq.common.message.*;\n\npublic class TransactionMessageService {\n \n private final TransactionMQProducer producer;\n \n public TransactionMessageService(String namesrvAddr) {\n this.producer = new TransactionMQProducer(\"tx-producer-group\");\n this.producer.setNamesrvAddr(namesrvAddr);\n this.producer.setTransactionListener(new OrderTransactionListener());\n this.producer.setCheckThreadPool(5);\n this.producer.setCheckThreadPoolMaxSize(10);\n }\n \n // 发送事务消息\n public TransactionSendResult sendOrderMessage(Order order) {\n Message msg = new Message(\n \"ORDER_TOPIC\", // topic\n \"OrderTag\", // tag\n order.getOrderId(), // keys(用于回查时定位本地事务)\n order.toJson().getBytes() // body\n );\n \n // 附加业务上下文,供本地事务使用\n msg.putUserProperty(\"bizType\", \"ORDER_CREATE\");\n msg.putUserProperty(\"amount\", String.valueOf(order.getAmount()));\n \n // 发送半消息(half message),消费者此时不可见\n TransactionSendResult result = producer.sendMessageInTransaction(msg, order);\n \n log.info(\"事务消息发送结果: txId={}, localState={}, sendStatus={}\",\n result.getTransactionId(),\n result.getLocalTransactionState(),\n result.getSendStatus());\n \n return result;\n }\n}\n\n// 事务监听器:实现本地事务 + 回查逻辑\nclass OrderTransactionListener implements TransactionListener {\n \n private final OrderService orderService;\n \n @Override\n public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {\n Order order = (Order) arg;\n try {\n // 1. 执行本地事务(创建订单)\n orderService.createOrder(order);\n // 2. 本地事务成功,提交半消息 → 消费者可见\n return LocalTransactionState.COMMIT_MESSAGE;\n } catch (DuplicateKeyException e) {\n // 3. 订单已存在(幂等),直接提交\n return LocalTransactionState.COMMIT_MESSAGE;\n } catch (Exception e) {\n // 4. 本地事务失败,回滚半消息 → 消费者不可见\n return LocalTransactionState.ROLLBACK_MESSAGE;\n }\n }\n \n @Override\n public LocalTransactionState checkLocalTransaction(MessageExt msg) {\n String orderId = msg.getKeys();\n \n // 回查逻辑:根据 orderId 查询本地事务状态\n Order order = orderService.queryOrder(orderId);\n \n if (order == null) {\n // 查不到 → 本地事务可能未执行或已回滚\n return LocalTransactionState.UNKNOW;\n }\n \n if (\"CREATED\".equals(order.getStatus()) || \"COMPLETED\".equals(order.getStatus())) {\n return LocalTransactionState.COMMIT_MESSAGE;\n }\n \n if (\"CANCELLED\".equals(order.getStatus())) {\n return LocalTransactionState.ROLLBACK_MESSAGE;\n }\n \n return LocalTransactionState.UNKNOW;\n }\n}",
"sub_questions": [
{
"index": 1,
"type": "single_choice",
"question": "当 executeLocalTransaction 抛出异常导致返回 ROLLBACK_MESSAGE 时,以下描述正确的是?",
"options": {
"A": "Broker 会将半消息标记为已回滚,消费者不会收到该消息,但消息仍会保留在 COMMIT_LOG 中",
"B": "Broker 会立即删除该半消息的物理存储",
"C": "Broker 会将该消息转发到死信队列供人工处理",
"D": "Broker 会自动重试 executeLocalTransaction 三次"
},
"answer": "A",
"explanation": "RocketMQ 事务消息回滚时,Broker 只是将半消息标记为回滚状态(在 HALF_TOPIC 中),使其对消费者不可见。消息的物理数据仍保留在 COMMIT_LOG 中(RocketMQ 采用追加写入,不支持物理删除),等待后续日志清理机制统一回收。选项 B 错误,RocketMQ 不做物理删除;选项 C 错误,回滚不是死信;选项 D 错误,回滚后不会重试本地事务。"
},
{
"index": 2,
"type": "short_answer",
"question": "当 checkLocalTransaction 返回 UNKNOW 时,RocketMQ Broker 的行为是什么?为什么要设计 UNKNOW 状态?",
"answer": "当 checkLocalTransaction 返回 UNKNOW 时,Broker 不会对半消息做任何操作(既不提交也不回滚),而是等待一段时间后再次发起回查。RocketMQ Broker 默认会以递增的时间间隔(如 60s, 120s, 240s...)最多回查 15 次。设计 UNKNOW 状态的原因是:本地事务可能正在执行中(如分布式事务协调器尚未完成),此时查询结果是不确定的,需要延迟再次确认。如果直接 COMMIT 或 ROLLBACK 可能导致数据不一致。",
"keywords": [
"UNKNOW",
"延迟回查",
"递增间隔",
"数据一致性",
"不确定性"
],
"scoring_rubric": "答出「Broker 会等待后再次回查」得 2 分;答出「递增间隔/多次回查」得 2 分;答出「UNKNOW 用于处理不确定状态/事务进行中」得 2 分。满分 6 分,4 分及以上通过。"
},
{
"index": 3,
"type": "single_choice",
"question": "代码中使用 order.getOrderId() 作为消息的 keys,这一设计对回查机制有什么关键作用?",
"options": {
"A": "仅用于消息检索,与回查机制无关",
"B": "回查时 Broker 将 keys 作为 MessageExt 的一部分下发,客户端通过 keys 定位对应的本地事务记录",
"C": "keys 会被 Broker 用于自动路由到正确的 Broker 节点",
"D": "keys 是消息的唯一标识,用于去重"
},
"answer": "B",
"explanation": "在 checkLocalTransaction 回调中,msg.getKeys() 返回的就是发送时设置的 keys 值(即 orderId)。回查的本质是 Broker 告诉 Producer「我有一个半消息不确定状态,请查一下你的本地事务」,Producer 需要一个标识来找到对应的本地记录。如果 keys 为空或不可唯一标识事务,回查时就无法判断本地事务是否成功。这正是代码中用 orderId 作为 keys 的核心原因。"
}
],
"explanation": "本题考查 RocketMQ 事务消息的完整生命周期:半消息发送 → executeLocalTransaction → 回查机制。事务消息的核心在于「本地事务 + 消息投递」的原子性保证。理解三种 LocalTransactionState(COMMIT/ROLLBACK/UNKNOW)的语义和 Broker 侧行为,以及 keys 在回查链路中的桥梁作用,是正确实现事务消息的关键。",
"source": null,
"related": []
},
{
"id": "cr-003",
"type": "code_reading",
"difficulty": 4,
"tags": [
"message-queue",
"dead-letter-queue",
"durable"
],
"question": "阅读以下 Go 语言实现的死信队列消费处理逻辑代码,分析其重试与兜底策略。",
"language": "go",
"code": "package mq\n\nimport (\n \"context\"\n \"encoding/json\"\n \"fmt\"\n \"log\"\n \"time\"\n)\n\n// DeadLetterHandler 处理死信队列中的消息\n\ntype DeadLetterHandler struct {\n retryProducer Producer\n alertService AlertService\n maxRetries int\n retryDelayBase time.Duration\n}\n\n// DeadLetterMessage 死信消息的结构\ntype DeadLetterMessage struct {\n OriginalTopic string `json:\"original_topic\"`\n OriginalBody []byte `json:\"original_body\"`\n ErrorInfo ErrorInfo `json:\"error_info\"`\n RetryCount int `json:\"retry_count\"`\n DLQTimestamp time.Time `json:\"dlq_timestamp\"`\n Metadata map[string]string `json:\"metadata\"`\n}\n\ntype ErrorInfo struct {\n ErrorCode string `json:\"error_code\"`\n Message string `json:\"message\"`\n}\n\n// HandleDeadLetter 处理单条死信消息\nfunc (h *DeadLetterHandler) HandleDeadLetter(ctx context.Context, msg *Message) error {\n var dlqMsg DeadLetterMessage\n if err := json.Unmarshal(msg.Body, &dlqMsg); err != nil {\n log.Printf(\"死信消息解析失败,丢弃: %v\", err)\n return nil // 解析失败的消息无法处理,确认消费避免无限循环\n }\n \n // 策略1:可重试的错误 → 重新投递到原 topic\n if h.isRetryable(dlqMsg.ErrorInfo.ErrorCode) {\n if dlqMsg.RetryCount < h.maxRetries {\n delay := h.calculateDelay(dlqMsg.RetryCount)\n \n err := h.retryProducer.SendWithDelay(\n ctx,\n dlqMsg.OriginalTopic,\n dlqMsg.OriginalBody,\n delay,\n )\n if err != nil {\n return fmt.Errorf(\"重试投递失败: %w\", err)\n }\n \n log.Printf(\"死信消息重试投递: topic=%s retry=%d/%d delay=%s\",\n dlqMsg.OriginalTopic, dlqMsg.RetryCount+1, h.maxRetries, delay)\n return nil\n }\n // 重试次数用尽,降级到告警\n log.Printf(\"重试次数用尽: topic=%s retry=%d\",\n dlqMsg.OriginalTopic, dlqMsg.RetryCount)\n }\n \n // 策略2:不可重试或重试用尽 → 告警 + 人工处理队列\n err := h.alertService.SendAlert(AlertPayload{\n Level: \"CRITICAL\",\n Title: \"死信消息需人工处理\",\n Detail: fmt.Sprintf(\"topic=%s error=%s retry=%d\",\n dlqMsg.OriginalTopic, dlqMsg.ErrorInfo.Message, dlqMsg.RetryCount),\n Timestamp: time.Now(),\n })\n if err != nil {\n log.Printf(\"告警发送失败: %v\", err)\n }\n \n // 将消息投递到人工处理队列\n return h.retryProducer.SendToQueue(ctx, \"MANUAL_DLQ_TOPIC\", msg.Body)\n}\n\nfunc (h *DeadLetterHandler) isRetryable(errorCode string) bool {\n retryableCodes := map[string]bool{\n \"NETWORK_TIMEOUT\": true,\n \"SERVICE_UNAVAILABLE\": true,\n \"RESOURCE_EXHAUSTED\": true,\n \"INVALID_INPUT\": false,\n \"PERMISSION_DENIED\": false,\n \"BUSINESS_RULE_VIOLATION\": false,\n }\n return retryableCodes[errorCode]\n}\n\nfunc (h *DeadLetterHandler) calculateDelay(retryCount int) time.Duration {\n // 指数退避:base * 2^retryCount,上限 30 分钟\n delay := h.retryDelayBase * time.Duration(1<<uint(retryCount))\n maxDelay := 30 * time.Minute\n if delay > maxDelay {\n delay = maxDelay\n }\n return delay\n}",
"sub_questions": [
{
"index": 1,
"type": "single_choice",
"question": "当死信消息解析失败(JSON Unmarshal 返回错误)时,代码选择返回 nil 而非 error,这样做的原因是?",
"options": {
"A": "解析失败的消息由上游负责重试,当前层不需要关心",
"B": "返回 error 会导致消息被重新投递到死信队列,形成无限循环",
"C": "解析失败意味着消息格式损坏,返回 error 会让 Broker 丢弃该消息",
"D": "这是一种错误处理的简化写法,功能上没有区别"
},
"answer": "B",
"explanation": "在消息队列的消费模型中,消费失败(返回 error)通常会导致消息被重试或重新投递到死信队列。对于一条格式已损坏的死信消息,如果返回 error,它会被再次投递到死信队列,再次解析失败,形成无限循环。返回 nil 意味着「我已成功处理(虽然实际上是放弃)」,确认消费后消息不会被重投。这是处理「无法挽救的消息」的标准防御性编程模式。"
},
{
"index": 2,
"type": "single_choice",
"question": "calculateDelay 使用指数退避算法 `base * 2^retryCount` 并设置 30 分钟上限,这样设计的主要目的是?",
"options": {
"A": "减少 Broker 的存储压力",
"B": "避免对下游服务产生突发重试风暴,同时保证问题恢复后消息仍能被及时处理",
"C": "节省 Producer 的网络带宽",
"D": "满足 Kafka 的消息保留策略要求"
},
"answer": "B",
"explanation": "指数退避(Exponential Backoff)的核心目的是在「给下游服务恢复时间」和「问题修复后及时重试」之间取得平衡。如果以固定间隔重试,当下游服务不可用时会产生大量无效请求(重试风暴);如果退避时间过长,服务恢复后消息处理会有不必要的延迟。设置 30 分钟上限确保即使重试多次,延迟也不会过长,保证业务时效性。"
},
{
"index": 3,
"type": "short_answer",
"question": "请分析该死信处理逻辑的三层兜底策略,并说明每层分别解决什么问题。",
"answer": "三层兜底策略:(1) 可重试错误 + 未超限 → 指数退避重投原 topic,解决临时性故障(网络超时、服务不可用等)导致的消息消费失败,给下游恢复机会;(2) 可重试但重试用尽 / 不可重试错误 → 发送告警通知运维人员,解决需要人工介入的故障场景(如业务规则冲突、权限问题等);(3) 最终兜底 → 投递到 MANUAL_DLQ_TOPIC 人工处理队列,确保消息不会丢失,为人工修复后重新处理提供入口。",
"keywords": [
"重试投递",
"告警",
"人工处理队列",
"指数退避",
"三层兜底",
"不可重试错误"
],
"scoring_rubric": "答出三层结构(重试/告警/人工队列)各得 2 分;能说明每层解决的问题类型(临时故障/人工介入/最终兜底)得 2 分。满分 8 分,5 分及以上通过。"
}
],
"explanation": "本题考查死信队列消费处理的核心设计模式。死信消息是消费失败的最终归宿,如何合理处理决定了系统的可靠性和可运维性。关键设计点包括:(1) 解析失败的防御性处理(避免无限循环);(2) 可重试/不可重试错误的分类策略;(3) 指数退避重试机制;(4) 告警与人工处理的最终兜底。这套模式在 Kafka、RocketMQ、RabbitMQ 等消息系统中普遍适用。",
"source": null,
"related": []
},
{
"id": "cr-004",
"type": "code_reading",
"difficulty": 3,
"tags": [
"message-queue",
"partition",
"cluster"
],
"question": "阅读以下 Java 代码,分析 Kafka 生产者自定义分区路由策略的实现逻辑。",
"language": "java",
"code": "package com.example.mq.partitioner;\n\nimport org.apache.kafka.clients.producer.Partitioner;\nimport org.apache.kafka.common.Cluster;\nimport org.apache.kafka.common.PartitionInfo;\nimport java.util.*;\nimport java.util.concurrent.ConcurrentHashMap;\n\n/**\n * 基于用户 ID 一致性哈希的分区策略\n * 保证同一用户的消息始终路由到同一分区,实现分区级顺序性\n */\npublic class ConsistentHashPartitioner implements Partitioner {\n \n private final Map<Integer, ConsistentHashRing> topicRings = new ConcurrentHashMap<>();\n private final int virtualNodes = 150; // 每个物理分区的虚拟节点数\n \n @Override\n public void configure(Map<String, ?> configs) {\n // 可从配置中读取虚拟节点数\n Object vnode = configs.get(\"consistent.hash.virtual.nodes\");\n if (vnode instanceof Integer) {\n // 此处省略赋值\n }\n }\n \n @Override\n public int partition(String topic, Object key, byte[] keyBytes,\n Object value, byte[] valueBytes, Cluster cluster) {\n List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);\n int numPartitions = partitions.size();\n \n if (key == null || keyBytes == null) {\n // 无 key 时使用默认轮询策略\n return defaultPartition(topic, numPartitions);\n }\n \n // 获取或创建该 topic 的哈希环\n ConsistentHashRing ring = topicRings.computeIfAbsent(topic,\n t -> buildHashRing(partitions));\n \n // 使用 murmur3 哈希确定路由\n int partition = ring.getNode(keyBytes);\n \n // 兜底:确保返回有效的分区编号\n return partition % numPartitions;\n }\n \n private ConsistentHashRing buildHashRing(List<PartitionInfo> partitions) {\n ConsistentHashRing ring = new ConsistentHashRing();\n for (PartitionInfo info : partitions) {\n ring.addPhysicalNode(info.partition(), virtualNodes);\n }\n return ring;\n }\n \n @Override\n public void close() {\n topicRings.clear();\n }\n \n private int defaultPartition(String topic, int numPartitions) {\n return Thread.currentThread().getId() % numPartitions;\n }\n}\n\n// 一致性哈希环实现\nclass ConsistentHashRing {\n private final TreeMap<Long, Integer> ring = new TreeMap<>();\n \n void addPhysicalNode(int nodeId, int virtualNodes) {\n for (int i = 0; i < virtualNodes; i++) {\n long hash = hash(nodeId + \"#\" + i);\n ring.put(hash, nodeId);\n }\n }\n \n int getNode(byte[] key) {\n long hash = hash(key);\n Map.Entry<Long, Integer> entry = ring.ceilingEntry(hash);\n if (entry == null) {\n entry = ring.firstEntry(); // 环形回绕\n }\n return entry.getValue();\n }\n \n private long hash(byte[] data) {\n // MurmurHash3 32-bit, 取绝对值确保非负\n return Math.abs(MurmurHash3.hash32(data));\n }\n}",
"sub_questions": [
{
"index": 1,
"type": "single_choice",
"question": "当 Broker 发生扩缩容(分区数变化)时,一致性哈希环相比普通取模分区的优势是什么?",
"options": {
"A": "只有约 1/N 的消息需要重新路由到新分区,N 为分区总数,减少消息重平衡范围",
"B": "扩缩容时不会有任何消息路由变化,完全无缝",
"C": "扩缩容时所有消息都会重新路由,但延迟更低",
"D": "一致性哈希环不支持分区数变化"
},
"answer": "A",
"explanation": "一致性哈希的核心优势在于:当节点数从 N 变为 N+1 时,虚拟节点环上的映射关系只影响约 1/(N+1) 的数据路由。相比普通取模分区(hash % N)在 N 变化时几乎所有 key 的分区都可能改变,一致性哈希大大减少了路由变化的范围,降低了扩缩容对消费者的影响(已消费的位点无需大量重算)。选项 B 错误,任何哈希策略在分区变化时都无法完全避免路由变化。"
},
{
"index": 2,
"type": "single_choice",
"question": "代码中为每个物理分区创建 150 个虚拟节点,虚拟节点的作用是什么?",
"options": {
"A": "增加网络连接数,提高消息吞吐量",
"B": "使消息在各物理分区间的分布更加均匀,减少数据倾斜",
"C": "提供消息备份,一个分区不可用时自动切换到虚拟节点",
"D": "用于实现消息的延迟投递"
},
"answer": "B",
"explanation": "物理分区数较少时(如 3-6 个),直接用物理节点构建哈希环会导致数据分布不均匀(某些分区承担过多路由)。虚拟节点通过为每个物理分区创建多个哈希映射点,使哈希环上的节点分布更密集、更均匀,从而让消息负载均衡到各个物理分区。150 个虚拟节点是常用的经验值,可以在均匀性和查找效率之间取得平衡。"
},
{
"index": 3,
"type": "short_answer",
"question": "代码中 defaultPartition 使用 Thread.currentThread().getId() % numPartitions 作为无 key 消息的分区策略,这种方式存在什么潜在问题?",
"answer": "这种无 key 时的分区策略存在以下问题:(1) 线程 ID 不一定均匀分布,可能导致分区倾斜——某些线程 ID 哈希后集中在少数分区;(2) 线程池大小变化(如扩容后线程数增加)会导致同一批消息的路由分布发生变化;(3) 与消费者组的线程模型耦合,如果消费者端线程数变化,生产端的分区分布也受影响。更稳健的做法是使用 AtomicInteger 的自增计数器取模,确保严格的轮询均匀分布。",
"keywords": [
"线程ID不均匀",
"分区倾斜",
"线程池变化",
"轮询",
"AtomicInteger"
],
"scoring_rubric": "答出「线程ID分布不均导致分区倾斜」得 3 分;答出「线程数变化影响路由分布」得 2 分;提出改进建议(如 AtomicInteger 轮询)得 1 分。满分 6 分,4 分及以上通过。"
}
],
"explanation": "本题考查 Kafka 自定义分区策略的设计与实现。Kafka 的默认分区策略(轮询或 key hash)在很多场景下不够灵活,通过实现 Partitioner 接口可以自定义路由逻辑。一致性哈希分区策略是保证「同一业务实体的消息顺序性」的经典方案,同时在扩缩容时具有良好的稳定性。理解虚拟节点、哈希环回绕、无 key 降级策略等细节是正确实现的前提。",
"source": null,
"related": []
},
{
"id": "cr-005",
"type": "code_reading",
"difficulty": 3,
"tags": [
"message-queue",
"consumer-group",
"durable"
],
"question": "阅读以下 Go 语言实现的消息幂等消费去重逻辑代码,分析其去重机制。",
"language": "go",
"code": "package mq\n\nimport (\n \"context\"\n \"crypto/sha256\"\n \"encoding/hex\"\n \"fmt\"\n \"log\"\n \"time\"\n)\n\n// IdempotentConsumer 幂等消费者\ntype IdempotentConsumer struct {\n store DeduplicationStore\n processor MessageProcessor\n windowSize time.Duration // 去重时间窗口\n}\n\ntype DeduplicationStore interface {\n // SetIfAbsent 原子操作:若 key 不存在则设置并返回 true,否则返回 false\n SetIfAbsent(ctx context.Context, key string, ttl time.Duration) (bool, error)\n // Get 获取去重记录的存在时间\n Get(ctx context.Context, key string) (time.Time, error)\n}\n\n// ConsumeMessage 幂等消费入口\nfunc (c *IdempotentConsumer) ConsumeMessage(ctx context.Context, msg *Message) error {\n // 1. 生成幂等键\n dedupKey := c.buildDedupKey(msg)\n \n // 2. 原子性检查 + 写入\n firstTime, err := c.store.SetIfAbsent(ctx, dedupKey, c.windowSize)\n if err != nil {\n // 存储层故障时的降级策略:放行,不做去重\n log.Printf(\"去重存储故障,降级放行: key=%s err=%v\", dedupKey, err)\n return c.processor.Process(ctx, msg)\n }\n \n if !firstTime {\n // 3. 重复消息,跳过处理\n log.Printf(\"重复消息已跳过: key=%s msgId=%s\", dedupKey, msg.ID)\n return nil\n }\n \n // 4. 首次消息,正常处理\n err = c.processor.Process(ctx, msg)\n if err != nil {\n // 5. 处理失败时需要删除去重键,允许重试\n c.store.Delete(ctx, dedupKey)\n return fmt.Errorf(\"消息处理失败,已清除去重键: %w\", err)\n }\n \n log.Printf(\"消息处理成功: key=%s msgId=%s\", dedupKey, msg.ID)\n return nil\n}\n\n// buildDedupKey 构造去重键:组合业务维度,避免跨业务误去重\nfunc (c *IdempotentConsumer) buildDedupKey(msg *Message) string {\n // 方案A:使用消息内置的唯一 ID(最简单)\n if msg.ID != \"\" {\n return fmt.Sprintf(\"dedup:%s\", msg.ID)\n }\n // 方案B:使用业务字段哈希(适用于无全局唯一 ID 的场景)\n raw := fmt.Sprintf(\"%s:%s:%s:%d\",\n msg.Headers[\"biz_type\"],\n msg.Headers[\"order_id\"],\n msg.Headers[\"action\"],\n msg.Timestamp.Unix(),\n )\n hash := sha256.Sum256([]byte(raw))\n return fmt.Sprintf(\"dedup:hash:%s\", hex.EncodeToString(hash[:8]))\n}",
"sub_questions": [
{
"index": 1,
"type": "single_choice",
"question": "当 SetIfAbsent 返回 false(非首次)时,代码直接返回 nil 不做任何处理。如果此时消费者需要「重新处理失败的消息」(即重试场景),这种设计会带来什么问题?",
"options": {
"A": "会导致消息被重复消费",
"B": "会导致重试的失败消息永远无法被重新处理,即「吞掉」重试消息",
"C": "没有任何问题,这是正确的去重逻辑",
"D": "会导致去重键永久占用存储空间"
},
"answer": "B",
"explanation": "这段代码的去重逻辑是「一次成功,永久跳过」。如果消息第一次处理失败(Process 返回 error),代码会删除去重键允许重试。但如果消息第一次处理成功后,由于某种原因需要重试处理(如业务层面的补偿),去重键仍然存在,重试消息会被当作重复消息跳过。这是去重逻辑的权衡——在「防重复」和「允许重试」之间需要根据业务场景选择不同的策略(如引入版本号或状态字段)。"
},
{
"index": 2,
"type": "single_choice",
"question": "当去重存储(Redis 等)故障时,代码选择「降级放行」而非「消费失败重试」,这样设计的理由是什么?",
"options": {
"A": "存储故障是永久性的,重试也没有意义",
"B": "保证消息处理的可用性优先于去重的严格性,避免因去重组件故障导致整体消费停滞",
"C": "Kafka 不支持消息重试",
"D": "降级放行的代码实现更简单"
},
"answer": "B",
"explanation": "分布式系统中,去重存储(如 Redis)是外部依赖,可能随时故障。如果此时消费失败重试,消息会被反复投递但始终无法处理,导致消费积压甚至消费者组崩溃。降级放行意味着「宁可重复消费,也不阻塞消费」,将可用性放在首位。重复消费的代价(如多扣一次款)可以通过业务层幂等(如数据库唯一约束)来兜底。这是分布式系统中「fail-open vs fail-close」的经典权衡。"
},
{
"index": 3,
"type": "short_answer",
"question": "代码中 buildDedupKey 提供了两种方案:消息 ID 和业务字段哈希。请分析各自的适用场景和优缺点。",
"answer": "方案 A(消息 ID):适用场景是消息系统提供了全局唯一 ID(如 Kafka 的 msg.Id 或 RocketMQ 的 MsgKey)。优点是简单可靠、无哈希冲突风险;缺点是依赖消息系统的 ID 唯一性保证,且不同消息系统或重试场景可能生成不同 ID,导致去重失效。\n方案 B(业务字段哈希):适用场景是消息没有全局唯一 ID,或需要按业务语义去重(如同一订单的同一操作)。优点是去重粒度可由业务控制,跨重试和消息系统保持一致;缺点是需要选择正确的业务字段组合,哈希存在理论上的碰撞风险,且时间精度(Unix秒级)可能影响去重精度。",
"keywords": [
"消息ID",
"业务字段哈希",
"全局唯一",
"去重粒度",
"哈希碰撞",
"适用场景"
],
"scoring_rubric": "方案 A 分析得当(适用场景 + 优缺点)得 3 分;方案 B 分析得当得 3 分;对比分析清晰得 1 分。满分 7 分,4 分及以上通过。"
}
],
"explanation": "本题考查消息幂等消费的核心设计模式。在分布式消息系统中,at-least-once 语义保证消息不丢但可能重复,幂等消费是处理重复消息的关键机制。核心设计点包括:(1) 原子性的去重检查(SetIfAbsent);(2) 去重键的设计(消息 ID vs 业务字段);(3) 存储故障的降级策略;(4) 处理失败时去重键的清理。在实际生产中,去重通常结合消息端(去重存储)和业务端(数据库唯一约束)双重保障。",
"source": null,
"related": []
}
]
}
@@ -0,0 +1,189 @@
{
"topic": "message-queue",
"type": "fill_blank",
"schema_version": "1.0.0",
"generated": "2026-09-09T16:12:31+08:00",
"questions": [
{
"id": "fb-001",
"type": "fill_blank",
"difficulty": 1,
"tags": [
"kafka",
"durable"
],
"question": "Kafka 中消息被追加到 Topic 的______末尾,这种写入方式被称为______写。",
"answer": [
"partition",
"append-only"
],
"answer_rule": "all",
"explanation": "Kafka 中 Topic 是逻辑概念,实际数据存储在 Partition(分区)中。每个 Partition 本质上是一个有序的、不可变的消息序列,新消息只能追加到末尾(append-only),这种顺序写方式极大减少了磁盘寻址开销,是 Kafka 高吞吐量的核心设计之一。",
"source": null,
"related": []
},
{
"id": "fb-002",
"type": "fill_blank",
"difficulty": 2,
"tags": [
"rocketmq",
"transactional-message"
],
"question": "RocketMQ 事务消息采用______消息机制:先发送半消息(Half Message)给 Broker,本地事务执行成功后再发送______(Commit/Rollback)指令,Broker 据此决定是否投递该消息。",
"answer": [
"two-phase",
"commit"
],
"answer_rule": "all",
"explanation": "RocketMQ 的事务消息基于两阶段提交思想。第一阶段发送 Half Message,此时消息对消费者不可见;第二阶段根据本地事务执行结果发送 Commit(提交,消息投递)或 Rollback(回滚,消息丢弃)。若 Broker 长时间未收到第二阶段指令,会主动回查生产者本地事务状态,确保最终一致性。",
"source": null,
"related": []
},
{
"id": "fb-003",
"type": "fill_blank",
"difficulty": 2,
"tags": [
"kafka",
"durable",
"zero-copy"
],
"question": "Kafka 利用操作系统的______系统调用实现零拷贝技术,将磁盘文件数据直接传输到网络 Socket,避免了用户态与内核态之间的数据拷贝。",
"answer": [
"sendfile"
],
"answer_rule": "all",
"explanation": "零拷贝(Zero-Copy)是 Kafka 高性能消费的关键技术。传统 I/O 需要经历:磁盘→内核缓冲区→用户缓冲区→内核 Socket 缓冲区→网卡,涉及多次拷贝和上下文切换。Kafka 通过 sendfile 系统调用,让数据直接从 PageCache 传输到网卡,仅需一次内核态操作,显著降低了消费延迟和 CPU 开销。",
"source": null,
"related": []
},
{
"id": "fb-004",
"type": "fill_blank",
"difficulty": 2,
"tags": [
"dead-letter-queue"
],
"question": "当消费者连续消费某条消息失败超过最大重试次数后,该消息会被投递到______队列,以避免阻塞后续正常消息的消费。",
"answer": [
"dead-letter",
"死信"
],
"answer_rule": "any",
"explanation": "死信队列(Dead Letter Queue,DLQ)是消息中间件中处理消费失败消息的重要机制。当消息消费失败次数超过配置的阈值(如 RocketMQ 默认 16 次重试),Broker 会将该消息转入死信 Topic。运维人员可以对死信队列中的消息进行人工排查、修复后重新投递或直接丢弃,防止问题消息无限重试阻塞消费流程。",
"source": null,
"related": []
},
{
"id": "fb-005",
"type": "fill_blank",
"difficulty": 3,
"tags": [
"kafka",
"transactional-message"
],
"question": "Kafka 事务机制中,生产者通过______ API 开启事务,所有跨分区的消息写入在该事务内是原子性的,即要么全部提交,要么全部回滚。",
"answer": [
"initTransactions / beginTransaction"
],
"answer_rule": "any",
"explanation": "Kafka 事务支持通过 Producer 的 initTransactions() 方法初始化事务协调器,然后在 sendMessages 过程中调用 beginTransaction() 开启事务,最后通过 commitTransaction() 或 abortTransaction() 提交或回滚。事务消息通过 Transaction Coordinator 和 __transaction_state 内部 Topic 来管理事务状态,保证跨分区的原子性写入,实现 Exactly-Once 语义。",
"source": null,
"related": []
},
{
"id": "fb-006",
"type": "fill_blank",
"difficulty": 3,
"tags": [
"kafka",
"cluster"
],
"question": "Kafka 集群中,每个 Partition 有一个 Leader 和多个______副本。只有 Leader 处理读写请求,Follower 仅负责同步数据。当 Leader 宕机时,Controller 从满足条件的 Follower 中选举新 Leader。",
"answer": [
"Follower",
"follower"
],
"answer_rule": "any",
"explanation": "Kafka 采用 Leader/Follower 主从架构。ISR(In-Sync Replicas)是与 Leader 保持同步的副本集合。Controller(集群中的一个 Broker)负责管理 Partition 的 Leader 选举。当 Leader 宕机,Controller 优先从 ISR 中选择第一个存活的 Follower 作为新 Leader,保证数据一致性。若 ISR 全部不可用,根据 unclean.leader.election.enable 配置决定是否允许不同步的副本接管。",
"source": null,
"related": []
},
{
"id": "fb-007",
"type": "fill_blank",
"difficulty": 3,
"tags": [
"push-pull",
"consumer-group"
],
"question": "Kafka 消费者采用______模式获取消息,消费者主动向 Broker 拉取数据。当某个消费者组中某个消费者宕机时,Broker 会触发______,将该消费者的 Partition 分配给组内其他消费者。",
"answer": [
"pull",
"rebalance"
],
"answer_rule": "all",
"explanation": "Kafka 采用 Pull(拉)模式,消费者根据自身处理能力主动从 Broker 拉取消息,相比 Push 模式有更好的流量控制能力。当 Consumer Group 中有消费者加入或离开时,会触发 Rebalance(重平衡),重新分配 Partition 的消费归属。Rebalance 期间消费暂停,因此需合理设置 session.timeout.ms 和 heartbeat.interval.ms 来平衡故障检测速度与误判风险。",
"source": null,
"related": []
},
{
"id": "fb-008",
"type": "fill_blank",
"difficulty": 4,
"tags": [
"rocketmq",
"kafka",
"topic"
],
"question": "RocketMQ 相比 Kafka,原生支持______机制,允许生产者在发送消息时附加 Tag,消费者可通过 Tag 过滤订阅感兴趣的消息,减少不必要的网络传输。",
"answer": [
"Tag",
"tag",
"消息过滤"
],
"answer_rule": "any",
"explanation": "RocketMQ 在 Topic 之下引入了 Tag 概念,生产者发送消息时可设置 Tag(如 OrderCreate、OrderCancel),消费者订阅时通过 Tag 表达式(如 * 或 || 或 &&)过滤消息。Broker 端支持基于 Tag 的服务端过滤,减少消费者拉取无关消息。Kafka 则依赖消息 Key 或消息体内容在客户端侧过滤,原生不支持服务端 Tag 过滤。",
"source": null,
"related": []
},
{
"id": "fb-009",
"type": "fill_blank",
"difficulty": 4,
"tags": [
"kafka",
"durable"
],
"question": "Kafka 的高性能写入依赖操作系统的______作为磁盘和网络之间的缓冲区。写入时数据先写入该缓存区,再由后台线程异步刷盘。若配置 acks=all 且______,则消息写入不会丢失。",
"answer": [
"PageCache",
"min.insync.replicas >= 2"
],
"answer_rule": "all",
"explanation": "Kafka 利用 PageCache(页缓存)实现高效的写入:生产者的消息先写入内核 PageCache,由 OS 后台线程异步刷盘(flush),避免同步 I/O 阻塞。配合 acks=all(所有 ISR 副本确认)和 min.insync.replicas >= 2 的配置,即使单个 Broker 宕机,只要 ISR 中还有其他副本持有数据,消息就不会丢失。这种设计在吞吐量和可靠性之间取得了平衡。",
"source": null,
"related": []
},
{
"id": "fb-010",
"type": "fill_blank",
"difficulty": 5,
"tags": [
"kafka",
"pulsar",
"transactional-message"
],
"question": "Kafka 的 Exactly-Once 语义通过幂等生产者(Producer)和______机制结合实现:幂等 Producer 为每条消息分配序列号,Broker 据此去重;事务机制则保证跨 Partition 写入的原子性。在消费端,需配合手动提交______和幂等消费逻辑才能端到端保证不丢不重。",
"answer": [
"transaction",
"offset"
],
"answer_rule": "all",
"explanation": "Kafka Exactly-Once 语义(EOS)由三层保障实现:(1) 幂等 Producer(enable.idempotence=true)通过 Producer ID + Sequence Number 在 Broker 端去重单分区内重复消息;(2) 事务机制(Transactional API)保证跨分区原子写入;(3) 消费端需配合 consumer.commitSync() 手动提交 offset,确保消费与提交在同一事务中完成(配合 transactional.id),再结合业务幂等(如数据库唯一键)实现端到端 Exactly-Once。",
"source": null,
"related": []
}
]
}
@@ -0,0 +1,25 @@
{
"slug": "message-queue",
"name": "消息队列",
"description": "主流中间件选型、路由、持久化原理、事务、死信队列、推拉模式、集群化部署",
"tags": [],
"question_files": {
"single_choice": "single_choice.json",
"true_false": "true_false.json",
"fill_blank": "fill_blank.json",
"short_answer": "short_answer.json",
"code_reading": "code_reading.json"
},
"stats": {
"total": 45,
"by_type": {
"single_choice": 15,
"true_false": 10,
"fill_blank": 10,
"short_answer": 5,
"code_reading": 5
}
},
"created": "2026-09-09",
"updated": "2026-09-09"
}
@@ -0,0 +1,148 @@
{
"topic": "message-queue",
"type": "short_answer",
"schema_version": "1.0.0",
"generated": "2026-09-09T15:30:00+08:00",
"questions": [
{
"id": "sa-001",
"type": "short_answer",
"difficulty": 2,
"tags": [
"kafka",
"cluster",
"durable"
],
"question": "简述 Kafka 中 ISR(In-Sync Replicas)机制的工作原理,以及它如何影响消息的可靠性保证。",
"answer": "ISR 是 Kafka 中与 Leader 保持同步的副本集合。当 Producer 发送消息时,通过 acks 参数控制可靠性级别:acks=0 不等待确认(可能丢消息);acks=1 等待 Leader 写入成功;acks=all/-1 等待所有 ISR 副本都写入成功后才返回确认。Follower 副本通过从 Leader 拉取数据保持同步,当 Follower 落后超过 replica.lag.time.max.ms 配置的时间时会被移出 ISR。Leader 选举时只会从 ISR 中选择新 Leader,从而保证已确认的消息不会丢失。ISR 机制在可靠性和可用性之间提供了灵活的平衡:ISR 越小,可用性越高但可靠性越低;ISR 越大,可靠性越高但写入延迟可能增加。此外,unclean.leader.election.enable 参数控制是否允许非 ISR 副本参与选举,进一步影响一致性与可用性的取舍。",
"keywords": [
"ISR",
"Leader",
"Follower",
"同步副本",
"acks",
"选举",
"可靠性和可用性权衡"
],
"scoring_rubric": "提到 ISR 定义(与 Leader 同步的副本集合)得 1 分;说明 acks 参数的三种级别及其影响得 1 分;解释 Follower 同步机制与 ISR 成员变化得 1 分;说明 ISR 与 Leader 选举的关系得 1 分;提及可靠性和可用性的权衡得 1 分。满分 5 分。",
"explanation": "ISR 是 Kafka 分布式消息系统中最核心的可靠性机制之一。理解 ISR 需要掌握:(1) Kafka 采用 Leader-Follower 副本模型,只有 Leader 处理读写;(2) Follower 通过拉取方式同步 Leader 数据;(3) 只有在 ISR 中的副本才有资格被选为新 Leader;(4) Producer 的 acks 配置决定了消息写入需要多少副本确认。这三个方面共同构成了 Kafka 的消息可靠性保证体系。面试中需要能清晰阐述各配置项的含义以及它们如何相互配合。",
"source": null,
"related": []
},
{
"id": "sa-002",
"type": "short_answer",
"difficulty": 5,
"tags": [
"rocketmq",
"transactional-message",
"exactly-once"
],
"question": "请详细解释 RocketMQ 事务消息的实现机制(半消息机制),并对比 Kafka 的事务消息方案,分析两者在实现 Exactly-Once 语义上的差异。",
"answer": "RocketMQ 事务消息采用半消息(Half Message)机制,核心流程分为四步:(1) Producer 先发送半消息到 Broker,消息对消费者不可见(处于 PREPARE 状态);(2) Producer 执行本地事务;(3) 根据本地事务执行结果,向 Broker 发送 COMMIT 或 ROLLBACK 命令;(4) 若 Broker 长时间未收到确认,会主动回查 Producer 的本地事务状态。事务消息存储在 Topic %RETRY% 下,通过回查机制(最多 15 次)保证最终一致性。Kafka 事务消息则基于 Transaction Coordinator 和 PID(Producer ID)机制:Producer 向 Coordinator 注册并获取 PID 和 Epoch,通过两阶段提交(AddPartitionsToTxn + EndTxn)实现事务。Kafka 事务是跨分区的,可以原子性地写入多个分区;RocketMQ 事务消息是单条消息级别的原子性。Exactly-Once 语义方面:RocketMQ 依赖事务消息+消费端幂等(MessageKey 去重)来近似实现;Kafka 通过幂等 Producer(PID+序列号去重)+ 事务消息的组合实现端到端 Exactly-Once,但消费者需要配合 read_committed 隔离级别。",
"keywords": [
"半消息",
"PREPARE",
"COMMIT",
"ROLLBACK",
"回查",
"Transaction Coordinator",
"PID",
"Exactly-Once",
"两阶段提交",
"幂等"
],
"scoring_rubric": "准确描述 RocketMQ 半消息的四步流程得 2 分;说明回查机制得 1 分;描述 Kafka 事务的 Coordinator+PID 机制得 2 分;对比两者在事务粒度上的差异(单条 vs 跨分区)得 1 分;分析 Exactly-Once 实现路径的差异得 1 分。满分 7 分。",
"explanation": "事务消息是消息队列中的高级话题,面试中常以对比形式考察。RocketMQ 的半消息机制是其独创设计,核心在于消息先投递但对消费者不可见,通过二次确认和回查机制保证事务的最终一致性。Kafka 的事务模型更偏向流处理场景,支持跨分区原子写入,但不提供消息回查机制。理解两者的差异需要掌握分布式事务的基本理论(两阶段提交、最终一致性),以及各自中间件的内部实现细节。这个题目难度较高,考察对底层机制的深入理解。",
"source": null,
"related": []
},
{
"id": "sa-003",
"type": "short_answer",
"difficulty": 3,
"tags": [
"dead-letter-queue",
"retry",
"message-trace"
],
"question": "什么是死信队列(Dead Letter Queue)?请说明死信产生的常见原因、典型的重试与处理策略,以及如何通过死信队列实现消息回溯。",
"answer": "死信队列(DLQ)是专门存放无法被正常消费的消息的特殊队列。死信产生的常见原因包括:(1) 消息格式非法或反序列化失败;(2) 消费逻辑抛出异常且达到最大重试次数;(3) 消息 TTL 过期;(4) 队列/Topic 满了导致消息被拒绝;(5) 消费者显式拒绝消息(如 RabbitMQ 的 basic.reject 或 basic.nack 且 requeue=false)。典型处理策略:(1) 退避重试:指数退避 + 最大重试次数,避免雪崩;(2) 死信队列消费:专门的消费者进程监听死信队列,进行人工审核或告警;(3) 消息修复与重放:修复后将消息重新投递到原队列;(4) 消息回溯:RocketMQ 支持按时间戳回溯消费(reconsume),Kafka 支持通过 --offset 重置消费者位点来回溯历史消息。消息追踪方面,可以通过在消息中嵌入 TraceID,配合链路追踪系统(如 SkyWalking、Jaeger)实现端到端的消息追踪和死信消息的全链路定位。",
"keywords": [
"死信队列",
"DLQ",
"重试",
"TTL",
"指数退避",
"回溯",
"TraceID",
"链路追踪"
],
"scoring_rubric": "准确定义死信队列得 1 分;列举至少三种死信产生原因得 1 分;说明退避重试策略得 1 分;描述死信队列的消费和处理方式得 1 分;说明消息回溯的具体方法得 1 分。满分 5 分。",
"explanation": "死信队列是消息系统可靠性设计的重要环节。生产环境中,消息消费失败是不可避免的,关键是如何优雅地处理这些失败消息。面试者需要理解死信队列不仅是'存放失败消息的地方',更是一套完整的异常消息处理体系,包括重试策略、告警机制、人工介入流程和消息回溯能力。理解 RabbitMQ、RocketMQ、Kafka 各自在死信处理上的差异也很重要。",
"source": null,
"related": []
},
{
"id": "sa-004",
"type": "short_answer",
"difficulty": 3,
"tags": [
"push-pull",
"consumer-group",
"rebalance"
],
"question": "对比消息队列中 Push 和 Pull 两种消费模式的优缺点,并说明 Long Polling 如何结合两者的优势。在消费者端流控和 Rebalance 方面,这两种模式分别面临哪些挑战?",
"answer": "Push 模式(Broker 主动推送给消费者):优点是实时性高、延迟低;缺点是 Broker 无法感知消费者处理能力,可能导致消费者过载(消息堆积)、背压控制困难。Pull 模式(消费者主动拉取消息):优点是消费者自主控制消费速率,天然支持流控和背压;缺点是存在轮询开销,实时性依赖拉取间隔。Long Polling(长轮询)结合两者优势:消费者发起拉取请求后,Broker 不立即返回空结果,而是等待有新消息或超时后才响应,既减少了无效轮询开销,又保证了实时性,同时保留了消费者端的流控能力。RabbitMQ 原生使用 Push 模式(basic.deliver),通过 prefetch 机制做流控;Kafka 和 RocketMQ 使用 Pull 模式,Kafka 的 high-level consumer 使用 long polling(fetch.min.bytes 配置等待时间)。Rebalance 方面:Push 模式下 Broker 需要维护每个消费者的状态,rebalance 时需要通知所有消费者;Pull 模式下消费者通过 consumer group 协议自行协调分区分配,rebalance 时会触发位点迁移,可能导致短暂的重复消费或消息延迟。Kafka 的 Rebalance 有 Eager(全停全分)和 Cooperative(增量分配)两种策略,后者可以减少 stop-the-world 影响。",
"keywords": [
"Push",
"Pull",
"Long Polling",
"prefetch",
"背压",
"流控",
"Rebalance",
"Eager",
"Cooperative",
"fetch.min.bytes"
],
"scoring_rubric": "准确对比 Push 和 Pull 的优缺点各得 1 分;说明 Long Polling 的结合优势得 1 分;讨论流控/背压挑战得 1 分;分析 Rebalance 相关挑战得 1 分。满分 5 分。",
"explanation": "推拉模式是消息队列的基础架构设计问题。面试中需要能清晰地对比两种模式的本质差异,理解 Long Polling 的工程价值,以及在生产环境中 Rebalance 带来的实际挑战(如消费暂停、重复消费、位点丢失等)。Kafka 的 Cooperative Rebalance 是近年的重要改进,了解这个演进能体现对技术发展的关注。",
"source": null,
"related": []
},
{
"id": "sa-005",
"type": "short_answer",
"difficulty": 4,
"tags": [
"kafka",
"rocketmq",
"pulsar",
"topic",
"partition",
"cluster"
],
"question": "从消息模型、持久化架构、运维复杂度、适用场景四个维度,对比 Kafka、RocketMQ 和 Apache Pulsar 三款消息中间件的异同,并给出各自的最佳适用场景。",
"answer": "消息模型:Kafka 采用 Topic-Partition 模型,消费者通过 offset 顺序消费,天然适合流式处理;RocketMQ 采用 Topic-Queue 模型,支持 Tag 和 Key 进行消息过滤和路由,更适合业务消息场景;Pulsar 采用 Topic-Partition + Segment 存储分离模型,支持多租户和租户级别的资源隔离,同时兼容 Kafka 的消费语义。持久化架构:Kafka 使用 Broker 本地磁盘顺序写 + PageCache + 零拷贝(sendfile),写入性能极高但存储与计算耦合;RocketMQ 类似,使用 CommitLog + ConsumeQueue 的双层存储结构,CommitLog 顺序写所有消息,ConsumeQueue 存储索引;Pulsar 采用 BookKeeper 存储层,实现了计算存储分离,写入 BookKeeper 后异步落盘,支持分层存储(Tiered Storage)。运维复杂度:Kafka 中等,依赖 ZooKeeper(新版本移除),分区再均衡需关注;RocketMQ 较简单,NameServer 无状态易部署,运维工具完善;Pulsar 较高,需要同时运维 Broker + BookKeeper + ZooKeeper 三个组件。适用场景:Kafka 适合大数据流处理、日志采集、实时数仓(高吞吐场景);RocketMQ 适合电商交易、金融核心链路、事务消息场景(可靠性和消息过滤需求);Pulsar 适合多租户 SaaS 平台、跨地域多活、需要存储计算分离的大规模消息平台。",
"keywords": [
"Topic",
"Partition",
"Queue",
"Tag",
"CommitLog",
"ConsumeQueue",
"BookKeeper",
"计算存储分离",
"零拷贝",
"多租户",
"Tiered Storage",
"ZooKeeper"
],
"scoring_rubric": "消息模型维度对比准确得 1 分;持久化架构维度对比准确得 1 分;运维复杂度维度对比准确得 1 分;适用场景划分合理得 1 分;总体表述清晰、有深度得 1 分。满分 5 分。",
"explanation": "三大消息中间件的对比是面试高频题。关键不在于死记参数,而在于理解设计哲学的差异:Kafka 以高吞吐和流处理为核心,追求极致的顺序写和零拷贝性能;RocketMQ 脱胎于电商场景,强调消息的可靠投递和灵活路由;Pulsar 则是新一代架构,通过计算存储分离解决 Kafka 的存储扩展性问题。面试者需要展现出对架构演进的理解,能够根据具体业务场景给出合理的选型建议。",
"source": null,
"related": []
}
]
}
@@ -0,0 +1,308 @@
{
"topic": "message-queue",
"type": "single_choice",
"schema_version": "1.0.0",
"generated": "2026-09-09T10:30:00+08:00",
"questions": [
{
"id": "sc-001",
"type": "single_choice",
"difficulty": 1,
"tags": [
"kafka",
"partition"
],
"question": "在 Kafka 中,一个 Topic 被分为多个 Partition,以下关于 Partition 的描述正确的是?",
"options": {
"A": "Partition 内的消息是全局有序的",
"B": "每个 Partition 内的消息是有序的,但不同 Partition 之间不保证顺序",
"C": "Partition 的数量在创建 Topic 时固定后不可更改",
"D": "每条消息会被广播到所有 Partition"
},
"answer": "B",
"explanation": "Kafka 的 Partition 是消息的物理分片单元,每个 Partition 内部的消息按照写入顺序严格有序(通过 offset 标识),但不同 Partition 之间不保证全局顺序。A 错误:全局有序需要只有一个 Partition。C 错误:Kafka 从 1.0 版本起支持动态增减 Partition 数量(但不能减少到低于当前最大 offset)。D 错误:消息通过分区器(Partitioner)根据 key 的 hash 值路由到特定的单个 Partition,而非广播。",
"source": null,
"related": []
},
{
"id": "sc-002",
"type": "single_choice",
"difficulty": 1,
"tags": [
"dead-letter-queue"
],
"question": "关于消息队列中的死信队列(Dead Letter Queue),以下说法正确的是?",
"options": {
"A": "死信队列中的消息会自动删除,无法再次消费",
"B": "死信队列用于存放消费失败且重试次数耗尽的消息",
"C": "死信队列只能由消息中间件自动创建,不支持手动配置",
"D": "死信队列中的消息优先级一定低于正常队列"
},
"answer": "B",
"explanation": "死信队列(DLQ)的核心作用是存放因消费失败(如格式错误、业务异常、消费超时等)且重试次数耗尽而无法正常消费的消息,便于后续人工排查或修复后重新投递。A 错误:DLQ 中的消息默认持久化保存,消费者可以再次消费或由运维人员处理。C 错误:大多数消息中间件(如 RocketMQ、RabbitMQ)都支持手动配置死信队列的名称和策略。D 错误:死信队列没有自动降低优先级的说法,它的优先级取决于消费者的消费顺序。",
"source": null,
"related": []
},
{
"id": "sc-003",
"type": "single_choice",
"difficulty": 2,
"tags": [
"kafka",
"rocketmq",
"pulsar"
],
"question": "以下关于 Kafka、RocketMQ 和 Pulsar 的对比,描述正确的是?",
"options": {
"A": "三者都采用 Broker + 存储耦合的架构设计",
"B": "Pulsar 采用计算与存储分离的架构,支持多租户和分层存储",
"C": "RocketMQ 的吞吐量一定高于 Kafka",
"D": "Kafka 原生支持事务消息,而 RocketMQ 不支持"
},
"answer": "B",
"explanation": "Pulsar 采用 BookKeeper 作为存储层、Broker 作为计算层的分离架构,天然支持多租户命名空间和分层存储(将冷数据卸载到 S3 等对象存储)。A 错误:Pulsar 是计算存储分离的,Kafka 和 RocketMQ 是存储耦合的。C 错误:Kafka 在高吞吐场景(日志收集、大数据流处理)中吞吐量通常优于 RocketMQ;RocketMQ 的优势在于低延迟和丰富的消息特性(如事务消息、延迟消息)。D 错误:Kafka 从 0.11 版本开始支持事务(Exactly-Once 语义),RocketMQ 也原生支持事务消息(半消息机制),两者都支持。",
"source": null,
"related": []
},
{
"id": "sc-004",
"type": "single_choice",
"difficulty": 2,
"tags": [
"consumer-group"
],
"question": "关于 Consumer Group 的行为,以下说法正确的是?",
"options": {
"A": "同一个 Consumer Group 内的不同消费者可以消费同一个 Partition 的消息",
"B": "不同 Consumer Group 之间消费进度互相影响",
"C": "同一个 Consumer Group 内,一条消息只会被其中一个消费者消费",
"D": "Consumer Group 的消费者数量必须等于 Partition 数量"
},
"answer": "C",
"explanation": "Consumer Group 是消息队列实现「负载均衡「和「广播「的核心机制。同一个 Group 内,一条消息只会被组内的一个消费者消费(点对点模式),实现负载均衡。A 错误:同一个 Group 内,一个 Partition 在同一时刻只能被一个消费者消费。B 错误:不同 Group 之间完全独立,各自维护独立的消费位点(offset),互不影响。D 错误:消费者数量可以少于或多于 Partition 数量;少于时一个消费者消费多个 Partition,多于时部分消费者空闲。",
"source": null,
"related": []
},
{
"id": "sc-005",
"type": "single_choice",
"difficulty": 2,
"tags": [
"push-pull"
],
"question": "关于消息队列的推(Push)拉(Pull)模式,以下描述正确的是?",
"options": {
"A": "Push 模式下消费者主动向 Broker 请求新消息",
"B": "Pull 模式下 Broker 主动将消息推送给消费者",
"C": "RocketMQ 采用 Push 模式时,底层实际上是基于长轮询(Long Polling)的 Pull 实现",
"D": "Kafka 的 Consumer 完全基于 Push 模式消费消息"
},
"answer": "C",
"explanation": "RocketMQ 的 Consumer 虽然名为 PushConsumer(API 表现为推送语义),但底层实现是基于长轮询(Long Polling)的 Pull 模式:消费者定期向 Broker 发起拉取请求,如果当前没有新消息,Broker 会 hold 住该请求(挂起一段时间),直到有新消息到达或超时后才返回。这种设计结合了 Pull 模式的消费者端流控能力和 Push 模式的消息及时性。A 错误:Push 模式是 Broker 主动推送。B 错误:Pull 模式是消费者主动拉取。D 错误:Kafka 的 Consumer 完全基于 Pull 模式,由消费者控制拉取频率。",
"source": null,
"related": []
},
{
"id": "sc-006",
"type": "single_choice",
"difficulty": 2,
"tags": [
"kafka",
"partition"
],
"question": "在 Kafka 中,以下哪种方式可以保证消息的全局顺序性?",
"options": {
"A": "将 Topic 的 Partition 数量设置为 3,并为消息设置 key",
"B": "将 Topic 的 Partition 数量设置为 1,所有消息不设置 key",
"C": "使用多个 Partition,通过消息 key 保证相同 key 的消息有序",
"D": "配置 acks=all 并设置 min.insync.replicas=2"
},
"answer": "B",
"explanation": "Kafka 只能保证单个 Partition 内的顺序性。要实现全局有序,必须将 Partition 数量设为 1,使所有消息写入同一个 Partition,从而利用 Partition 内的 offset 顺序保证全局有序。A 错误:3 个 Partition 只能保证每个 Partition 内有序,Partition 之间不保证顺序。C 错误:相同 key 的消息会被路由到同一个 Partition,保证了 key 级别的局部有序,但不同 key 的消息仍然跨 Partition 分布,不保证全局有序。D 错误:acks 和 min.insync.replicas 控制的是数据可靠性和一致性,与消息顺序性无关。",
"source": null,
"related": []
},
{
"id": "sc-007",
"type": "single_choice",
"difficulty": 3,
"tags": [
"kafka",
"durable",
"topic"
],
"question": "Kafka 的高吞吐量很大程度上依赖于其对零拷贝(Zero-Copy)技术的使用,关于该技术以下描述正确的是?",
"options": {
"A": "零拷贝通过内核态和用户态之间的内存映射实现数据传输",
"B": "零拷贝使用 sendfile 系统调用,使数据直接从磁盘(PageCache)传输到网卡,绕过用户空间",
"C": "零拷贝只能用于消费者首次消费消息的场景",
"D": "零拷贝要求消息必须存储在内存中而非磁盘上"
},
"answer": "B",
"explanation": "Kafka 的 Consumer Fetch 数据时,使用 Linux 的 sendfile 系统调用实现零拷贝。传统 I/O 需要经过「磁盘→内核缓冲区→用户缓冲区→Socket 缓冲区→网卡」四次拷贝,而 sendfile 直接在内核态将数据从 PageCache 传输到网卡(DMA gather copy),绕过了用户空间,减少了上下文切换和内存拷贝的开销。A 错误:零拷贝的核心是 sendfile,不是 mmap。C 错误:零拷贝对所有消费读取都适用(包括重复消费和多消费者场景)。D 错误:零拷贝读取的数据来自 PageCache(磁盘数据在内核页缓存中的映射),不要求数据在应用层内存中。",
"source": null,
"related": []
},
{
"id": "sc-008",
"type": "single_choice",
"difficulty": 3,
"tags": [
"kafka",
"cluster"
],
"question": "关于 Kafka 的 ISR(In-Sync Replicas)机制,以下描述正确的是?",
"options": {
"A": "ISR 是所有副本的集合,包括与 Leader 完全同步的和落后的副本",
"B": "ISR 中的副本必须与 Leader 保持数据同步,落后太多的副本会被移出 ISR",
"C": "当 ISR 中所有副本都确认写入后,Producer 才会收到 ack,这会导致可用性降低",
"D": "ISR 的大小永远等于 Replication Factor 的值"
},
"answer": "B",
"explanation": "ISR 是与 Leader 保持数据同步的副本集合。Kafka 通过 replica.lag.time.max.ms 参数控制副本的最大允许落后时间,如果某个 Follower 副本落后超过该阈值,会被移出 ISR。这样保证了 ISR 中的副本都有能力在 Leader 故障时被选为新 Leader。A 错误:ISR 严格排除了落后过多的副本。C 错误:acks=all 只要求 ISR 中所有副本确认,ISR 缩小反而降低了写入延迟;可用性降低主要在 ISR 缩为 1 且该节点故障时发生。D 错误:ISR 大小是动态变化的,可能小于 Replication Factor(因副本落后被踢出),也可能等于(正常情况)。",
"source": null,
"related": []
},
{
"id": "sc-009",
"type": "single_choice",
"difficulty": 3,
"tags": [
"rocketmq",
"transactional-message"
],
"question": "关于 RocketMQ 的事务消息机制,以下描述正确的是?",
"options": {
"A": "事务消息的半消息(Half Message)在发送后消费者立即可见",
"B": "事务消息通过两阶段提交实现:先发半消息,本地事务执行成功后发送 Commit,失败则发送 Rollback",
"C": "RocketMQ 事务消息依赖数据库 XA 事务来保证一致性",
"D": "事务消息的回查机制是由消费者端触发的"
},
"answer": "B",
"explanation": "RocketMQ 事务消息采用两阶段提交:Producer 先发送半消息(Half Message)到 Broker,此时消费者不可见;然后 Producer 执行本地事务(如数据库操作);成功则发送 Commit 使消息对消费者可见,失败则发送 Rollback 删除消息。A 错误:半消息在 Commit 之前对消费者不可见,这是事务消息隔离性的关键。C 错误:RocketMQ 事务消息不依赖数据库 XA 事务,而是通过消息回查(Transaction Check)机制——如果 Broker 未收到 Commit 或 Rollback,会定期回查 Producer 本地事务状态来决定最终提交或回滚。D 错误:回查机制是由 Broker 端主动发起,查询 Producer 的本地事务状态。",
"source": null,
"related": []
},
{
"id": "sc-010",
"type": "single_choice",
"difficulty": 3,
"tags": [
"kafka",
"durable"
],
"question": "Kafka 顺序写磁盘的性能优势来源于什么?",
"options": {
"A": "顺序写避免了磁盘寻道时间,同时利用了操作系统的 PageCache 做读写加速",
"B": "顺序写时磁盘的 IOPS 指标会显著提升",
"C": "顺序写意味着每条消息只写入一次,不需要追加写",
"D": "顺序写通过 RAID 0 条带化实现并行写入多个磁盘"
},
"answer": "A",
"explanation": "Kafka 的核心设计哲学是将磁盘当作顺序写入的存储来使用。顺序写磁盘避开了随机 I/O 的磁盘寻道开销(机械硬盘的寻道时间通常 5-10ms),在顺序写入场景下磁盘吞吐量可达 600MB/s 以上。同时,Kafka 利用操作系统的 PageCache 作为读写缓冲层:写入时数据先写入 PageCache(用户态到内核态的一次拷贝),再由 OS 异步刷盘;读取时优先从 PageCache 返回,未命中时再读磁盘。B 错误:IOPS 是衡量随机 I/O 的指标,顺序写关注的是吞吐量(Throughput)。C 错误:Kafka 是 append-only 的追加写模式,每条消息确实只写一次,但这不是顺序写的性能来源。D 错误:RAID 条带化是硬件层面的并行方案,不是 Kafka 顺序写的设计原理。",
"source": null,
"related": []
},
{
"id": "sc-011",
"type": "single_choice",
"difficulty": 3,
"tags": [
"kafka",
"consumer-group",
"push-pull"
],
"question": "当 Consumer Group 内有消费者加入或离开时,Kafka 会触发 Rebalance。关于 Rebalance,以下描述正确的是?",
"options": {
"A": "Rebalance 过程中消费者可以继续正常消费消息",
"B": "Rebalance 由 Consumer 端主动发起,Broker 不参与分配决策",
"C": "Rebalance 会导致所有消费者暂停消费,直到新的分区分配方案确定",
"D": "Rebalance 只在消费者宕机时才会触发"
},
"answer": "C",
"explanation": "Kafka 的 Rebalance 是将 Partition 重新分配给 Consumer Group 内的消费者的过程。在 Rebalance 期间,所有消费者会暂停消费(进入 rebalancing 状态),直到分配方案确定后才恢复。A 错误:Rebalance 期间消费者无法消费消息,这是一个已知的延迟来源。B 错误:在 Kafka 2.3+ 的 Cooperative Rebalance(增量再平衡)中 Broker(Group Coordinator)参与协调,且新版 Kafka 支持 Sticky 分配策略,由协调者参与决策。D 错误:Rebalance 不仅在宕机时触发,消费者主动加入/离开(如发布新版本重启、扩容缩容)、心跳超时(session.timeout.ms)、Poll 超时(max.poll.interval.ms)等都会触发 Rebalance。",
"source": null,
"related": []
},
{
"id": "sc-012",
"type": "single_choice",
"difficulty": 4,
"tags": [
"kafka",
"transactional-message"
],
"question": "关于 Kafka 的 Exactly-Once 语义实现,以下描述正确的是?",
"options": {
"A": "Kafka 通过acks=all就能实现端到端的Exactly-Once语义",
"B": "Kafka的Exactly-Once语义依赖于幂等Producer(PID+序列号去重)和事务API的组合",
"C": "Kafka的Exactly-Once语义仅适用于Producer端,Consumer端无法实现",
"D": "Kafka的Exactly-Once语义要求所有Partition的Replication Factor必须为1"
},
"answer": "B",
"explanation": "Kafka 的 Exactly-Once 语义由两个机制组合实现:(1) 幂等 Producer——Broker 为每个 Producer 分配唯一的 PID(Producer ID),每条消息携带递增的序列号,Broker 端通过 PID+序列号去重,解决单个 Producer 到 Broker 的重复问题;(2) 事务 API——Producer 可以在一个原子事务中写入多条消息(跨 Partition),确保要么全部提交要么全部回滚,解决跨 Partition 的原子写入问题。A 错误:acks=all 只保证消息被所有 ISR 副本接收,不能解决 Producer 重试导致的重复发送问题,更不是端到端的 Exactly-Once。C 错误:Consumer 配合 read_committed 隔离级别和幂等消费逻辑(如数据库唯一键),也可以实现端到端的 Exactly-Once。D 错误:Replication Factor 与 Exactly-Once 无关。",
"source": null,
"related": []
},
{
"id": "sc-013",
"type": "single_choice",
"difficulty": 4,
"tags": [
"kafka",
"durable"
],
"question": "以下关于 Kafka 中 PageCache 和 fsync 策略的描述,哪项是正确的?",
"options": {
"A": "Kafka 默认使用同步刷盘(每条消息写入后立即 fsync),以确保数据不丢失",
"B": "Kafka 依赖操作系统的 PageCache 做写缓冲,由操作系统异步刷盘,这在 broker 崩溃时可能导致少量数据丢失",
"C": "配置 log.flush.interval.messages=1 可以实现每条消息都 fsync,同时不影响写入吞吐量",
"D": "PageCache 是用户态的内存缓存,用于缓存 Kafka 的日志数据"
},
"answer": "B",
"explanation": "Kafka 的默认设计是利用操作系统的 PageCache 作为写入缓冲:Producer 写入的数据先到达 PageCache(内核态),然后由操作系统的 pdflush 等机制异步刷盘。这意味着如果 Broker 进程崩溃(OS 仍正常运行),PageCache 中未刷盘的数据不会丢失;但如果 OS 也崩溃(断电),未刷盘的数据可能丢失。这是 Kafka 在吞吐量和可靠性之间的权衡。A 错误:Kafka 默认不是同步刷盘,否则吞吐量会大幅下降。C 错误:虽然配置 log.flush.interval.messages=1 可以每条消息 fsync,但这会严重降低写入吞吐量,是已知的性能反模式。D 错误:PageCache 是内核态的内存页缓存,不是用户态的。",
"source": null,
"related": []
},
{
"id": "sc-014",
"type": "single_choice",
"difficulty": 4,
"tags": [
"dead-letter-queue",
"rocketmq"
],
"question": "在 RocketMQ 中,关于死信队列和消息重试的机制,以下描述正确的是?",
"options": {
"A": "消息消费失败后直接进入死信队列,不进行任何重试",
"B": "RocketMQ 的重试队列支持按时间延迟重试,重试次数耗尽后消息自动转入死信队列",
"C": "死信队列中的消息只能由运维人员手动删除,不支持再次消费",
"D": "RocketMQ 的重试次数由 Broker 全局统一配置,消费者无法自定义"
},
"answer": "B",
"explanation": "RocketMQ 的消息重试机制:当消息消费失败时,Broker 会将消息投递到内置的重试队列(%RETRY%{ConsumerGroup})。重试队列按照延迟级别组织(默认 18 个级别:10s, 30s, 1m, 2m, 3m, ... 2h),消息会按照延迟级别逐步重试。当重试次数超过配置的最大重试次数(默认 16 次)后,消息会被自动转入死信队列(%DLQ%{ConsumerGroup})。A 错误:消费失败后会先进入重试队列进行延迟重试,而非直接进死信队列。C 错误:死信队列中的消息可以被消费者消费(通常是由专门的死信消费者来处理),也可以选择删除。D 错误:最大重试次数由消息的属性(reconsumeTimes)控制,消费者可以通过修改消息属性来自定义重试行为。",
"source": null,
"related": []
},
{
"id": "sc-015",
"type": "single_choice",
"difficulty": 5,
"tags": [
"kafka",
"cluster"
],
"question": "在 Kafka 集群中,当 Leader 副本所在的 Broker 宕机后,以下关于 Leader 选举的描述正确的是?",
"options": {
"A": "Kafka 使用 Paxos 算法进行 Leader 选举,保证强一致性",
"B": "ISR 中的存活副本会参与 Leader 选举,由 Controller 从 ISR 中按优先级选择新 Leader",
"C": "所有 Follower 副本都会参与投票,得票最多的成为新 Leader",
"D": "Leader 选举完成后,所有 Follower 需要从新的 Leader 全量同步数据才能开始提供读服务"
},
"answer": "B",
"explanation": "Kafka 的 Leader 选举由集群中的 Controller(一个特殊的 Broker,通过 ZooKeeper/KRaft 选举产生)负责管理。当 Leader 副本故障时,Controller 从该 Partition 的 ISR 列表中选择一个存活的副本作为新 Leader(如果配置了 unclean.leader.election.enable=true,ISR 全部不可用时还可能从非 ISR 副本中选择,但会有数据丢失风险)。新 Leader 选定后立即对外提供服务,其他 Follower 在新 Leader 上开始拉取数据逐步同步。A 错误:Kafka 不使用 Paxos,早期使用 ZooKeeper 的 ZAB 协议,新版 Kafka 使用内置的 KRaft(基于 Raft 的变体)进行元数据管理。C 错误:Kafka 不采用投票机制,而是由 Controller 直接从 ISR 中选择。D 错误:新 Leader 选出后立即可用,Follower 异步追赶数据,不需要全量同步完成后才提供服务(这正是 ISR 机制的意义——ISR 中的副本数据已经足够接近 Leader)。",
"source": null,
"related": []
}
]
}
@@ -0,0 +1,160 @@
{
"topic": "message-queue",
"type": "true_false",
"schema_version": "1.0.0",
"generated": "2026-09-09T15:30:00+08:00",
"questions": [
{
"id": "tf-001",
"type": "true_false",
"difficulty": 3,
"tags": [
"kafka",
"零拷贝",
"持久化"
],
"question": "Kafka 的零拷贝(Zero-Copy)技术依赖于 Linux 操作系统的 sendfile 系统调用,Java 层面通过 FileChannel.transferTo() 方法触发,其核心优势是避免了数据在用户态和内核态之间的拷贝。",
"answer": true,
"explanation": "正确。Kafka 在将磁盘文件数据发送到网络 Socket 时,使用了操作系统的 sendfile 系统调用(通过 Java NIO 的 FileChannel.transferTo() 触发)。sendfile 直接在内核空间完成文件数据到网卡缓冲区的传输,无需将数据拷贝到用户态再写回内核态,显著减少了 CPU 开销和内存拷贝次数,是 Kafka 高吞吐量的关键优化之一。",
"source": null,
"related": []
},
{
"id": "tf-002",
"type": "true_false",
"difficulty": 2,
"tags": [
"rocketmq",
"事务消息",
"半消息"
],
"question": "RocketMQ 的事务消息机制依赖于半消息(Half Message)实现:Producer 先发送半消息到 Broker,执行本地事务后再根据结果 commit 或 rollback,Broker 对消费者隐藏未确认的半消息。",
"answer": true,
"explanation": "正确。这是 RocketMQ 事务消息的核心机制。Producer 发送半消息后,Broker 会将其存储在内部的 Half Topic 中,此时消费者无法消费该消息。Producer 执行本地事务后,根据结果向 Broker 发送 commit 或 rollback 指令。commit 后消息进入真正的 Topic 可被消费;rollback 后消息被丢弃。若 Broker 未收到确认,会定时回查 Producer 的本地事务状态。",
"source": null,
"related": []
},
{
"id": "tf-003",
"type": "true_false",
"difficulty": 4,
"tags": [
"pulsar",
"架构",
"计算存储分离"
],
"question": "Apache Pulsar 采用计算存储分离架构,Broker 节点是无状态的,消息数据持久化在 Apache BookKeeper 中;而 Kafka 的 Broker 节点同时承担计算和存储职责,消息数据直接存储在 Broker 本地磁盘上。",
"answer": true,
"explanation": "正确。这是 Pulsar 和 Kafka 架构的根本区别。Pulsar 的 Broker 不存储数据,仅负责协议处理和读写路由,实际数据写入 BookKeeper 的 Bookie 节点;而 Kafka 的 Broker 既处理客户端请求(计算),又将数据存储在本地磁盘的日志段中(存储)。这种分离使得 Pulsar 在 Broker 扩缩容时无需数据迁移,而 Kafka 扩容 Broker 后需要进行分区重分配(rebalance)。",
"source": null,
"related": []
},
{
"id": "tf-004",
"type": "true_false",
"difficulty": 2,
"tags": [
"rocketmq",
"推拉模式",
"长轮询"
],
"question": "RocketMQ 的 Consumer 采用 Push(推)模式消费消息,即 Broker 主动将消息推送给消费者;而 Kafka 的 Consumer 始终采用 Pull(拉)模式,由消费者主动向 Broker 拉取消息。",
"answer": false,
"explanation": "错误。虽然 RocketMQ 的 Consumer API 是 Push 模式(注册 MessageListener 回调),但底层实现实际上是基于长轮询(Long Polling)的 Pull 模式。Consumer 向 Broker 发起长轮询请求,Broker 在有新消息时立即响应,无消息时保持连接等待。Kafka 的 Consumer 确实是 Pull 模式。因此 RocketMQ 并非真正的 Broker 主动推送,而是通过长轮询模拟了 Push 的实时性。",
"source": null,
"related": []
},
{
"id": "tf-005",
"type": "true_false",
"difficulty": 3,
"tags": [
"kafka",
"集群",
"ISR",
"选举"
],
"question": "Kafka 的 ISR(In-Sync Replicas)机制仅用于控制 follower 副本从 leader 副本同步数据的时延阈值(replica.lag.time.max.ms),与数据一致性和容灾选举无关。",
"answer": false,
"explanation": "错误。ISR 是 Kafka 保障数据一致性和可用性的核心机制,绝不仅仅控制同步延迟。ISR 维护的是与 Leader 保持同步的副本集合:当 follower 副本落后超过 replica.lag.time.max.ms 时会被踢出 ISR。ISR 的作用包括:1)acks=all 时只有 ISR 中所有副本确认才算写入成功,保障数据一致性;2)Leader 宕机时,Controller 优先从 ISR 中选举新 Leader,保障可用性和数据完整性。ISR 是数据一致性和容灾选举的关键枢纽。",
"source": null,
"related": []
},
{
"id": "tf-006",
"type": "true_false",
"difficulty": 3,
"tags": [
"kafka",
"事务消息",
"exactly-once"
],
"question": "Kafka 的事务消息机制允许 Producer 在一个事务内跨多个分区原子性地写入消息,配合 Consumer 的 isolation.level=read_committed 配置,可实现端到端的 Exactly-Once 语义。",
"answer": true,
"explanation": "正确。Kafka 从 0.11 版本引入事务支持。Producer 通过 initTransactions()、beginTransaction()、commitTransaction() 等 API 在一个事务内跨分区写入消息,由 TxnCoordinator 协调事务状态。Consumer 设置 isolation.level=read_committed 后,只能读取已提交事务的消息。结合幂等 Producer(enable.idempotence=true),Kafka 可在 Producer→Broker→Consumer 链路上实现 Exactly-Once 语义。",
"source": null,
"related": []
},
{
"id": "tf-007",
"type": "true_false",
"difficulty": 2,
"tags": [
"消息模型",
"广播",
"consumer-group"
],
"question": "在 RocketMQ 中,广播消费(BROADCASTING)模式下,同一个 Consumer Group 中的每个 Consumer 实例都会收到全量消息的副本;而集群消费(CLUSTERING)模式下,同组内只有一个实例收到某条消息。",
"answer": true,
"explanation": "正确。RocketMQ 支持两种消费模式:广播消费(BROADCASTING)下,消息会投递给 Consumer Group 中的每一个实例,每个实例独立消费全量消息,适用于通知、配置更新等场景;集群消费(CLUSTERING)下,同组内的消息被负载均衡分配给其中一个实例处理,适用于任务分发场景。注意:广播消费模式下 rebalance 不生效,因为每个实例都需要消费所有消息。",
"source": null,
"related": []
},
{
"id": "tf-008",
"type": "true_false",
"difficulty": 4,
"tags": [
"死信队列",
"重试策略",
"消息可靠性"
],
"question": "在 RocketMQ 中,当消费端处理消息失败后,消息会自动进入重试队列进行多次重试;超过最大重试次数后,消息会被自动路由到死信队列(DLQ),无需开发者额外配置死信队列的消费逻辑即可完成消息回溯。",
"answer": false,
"explanation": "错误。前半句正确——RocketMQ 确实会在消费失败后将消息放入重试 Topic(%RETRY%)进行重试,达到最大重试次数后自动转入死信队列(%DLQ%)。但后半句错误:死信队列中的消息不会自动被消费或回溯。开发者必须编写专门的消费者订阅死信队列 Topic 进行处理(如告警、人工干预或写入落库)。如果不对死信队列进行消费,消息将永久滞留在 DLQ 中,造成消息丢失。",
"source": null,
"related": []
},
{
"id": "tf-009",
"type": "true_false",
"difficulty": 5,
"tags": [
"kafka",
"消费者",
"offset",
"exactly-once"
],
"question": "Kafka 消费者在处理完消息后立即提交 offset(即 process-then-commit 模式),配合自动提交机制,即可实现 Exactly-Once 消费语义,无需额外的幂等处理。",
"answer": false,
"explanation": "错误。这里存在两个陷阱:1)「自动提交」(enable.auto.commit=true)是在 poll() 调用时自动提交上次的 offset,而非在消息处理完成后提交,这意味着如果消费者在处理消息过程中崩溃,已提交的 offset 对应的消息实际上未被处理,导致消息丢失(At-Most-Once);2)即使手动在处理完成后提交 offset(At-Least-Once),也只是保证消息至少被处理一次,仍可能出现重复消费。要实现 Exactly-Once,需要结合幂等消费(如数据库唯一键去重)或使用 Kafka 事务将 offset 提交与消息处理绑定在同一事务中。process-then-commit + 自动提交不能实现 Exactly-Once。",
"source": null,
"related": []
},
{
"id": "tf-010",
"type": "true_false",
"difficulty": 2,
"tags": [
"pulsar",
"消息模型",
"确认机制"
],
"question": "Pulsar 同时支持累积确认(Cumulative Acknowledgment)和单条确认(Individual Acknowledgment)两种消费确认模式,而 Kafka Consumer 仅支持对 poll() 批次内已处理消息的批量提交,不支持选择性确认单条消息。",
"answer": false,
"explanation": "错误。后半句关于 Kafka 的描述不准确。Kafka 从 0.11 版本起支持通过 ConsumerRecord 记录每条消息的 offset,开发者可以使用 consumer.commitSync(Collection<ConsumerPartitionOffsetMetadata>) 方法选择性地提交特定分区和 offset 的确认,实现类似「单条确认」的效果。此外,Kafka 还支持通过 ConsumerRebalanceListener 在 rebalance 时精确控制 offset 提交。Pulsar 确实原生支持累积确认和单条确认两种模式,但说 Kafka 不支持选择性确认是错误的。",
"source": null,
"related": []
}
]
}