代表性面试主题

Kafka 后端面试:Share Group 如何提供队列语义?

后端困难
Offer.cc 编辑团队发布 更新

题干

一个订单处理服务希望让多个消费者共同领取同一 Kafka topic 的独立任务,并支持逐条确认和失败重试。请比较 Share Group 与 Consumer Group,说明何时可以采用、如何处理顺序与重复,以及如何安全迁移。

题干与适用场景

订单通知、图片转码或账单计算通常是彼此独立的工作项。团队已经把事件写入 Kafka,但传统 Consumer Group 让每个 partition 在同一时刻只分配给一个消费者;当消费者数量超过 partition 数量时,额外实例不能直接增加并行度。请设计一个方案,让多个消费者协作领取工作,支持逐条确认、失败重试和可观测的投递次数,同时明确哪些业务仍应保留 partition 顺序。

这里讨论的是 Apache Kafka KIP-932 定义的 Share Group。KIP 将它描述为在普通 topic 上提供协作式消费的 group 类型,并不等于把 Kafka 变成具有相同语义的 RabbitMQ。回答必须先确认部署版本、客户端支持和相关 API 是否已允许使用,再决定是否上线。

面试官考察点

  • 能否解释 partition 独占分配与 Share Group 协作领取的差异。
  • 能否把“确认、释放、拒绝、锁过期”映射到实际处理状态。
  • 是否知道 Share Group 允许消费者数超过 partition 数,但不自动保留传统 key 顺序。
  • 是否把重复投递、毒性消息、处理超时和最大并发作为同一套失败模型讨论。
  • 是否会核对客户端、broker、权限、监控和回滚,而不是只背 KIP 名称。

普通回答会说“Kafka 也能当队列”;强回答会指出队列语义的收益、失去的保证和上线前必须验证的边界。

回答前需要澄清的问题

  1. 工作项是否真正独立?若同一订单的事件必须按 key 顺序处理,Share Group 可能改变解法。
  2. 失败时是短暂重试、转入人工处理,还是永久拒绝?这决定 release、reject 和死信流程。
  3. 处理时长的 p99、最大并发和可接受的重复副作用是多少?它们决定 acquisition lock 与幂等设计。
  4. 业务需要 Kafka 事务的端到端语义吗?不能假设 Share Group 自动继承既有 Consumer Group 事务方案。
  5. 现有 topic 是否还服务广播或回放消费者?迁移 group 类型不能改变其他 group 的读取契约。

30 秒回答框架

“我先确认任务是否允许乱序,以及部署的 Kafka 和客户端是否支持 KIP-932。Share Group 让多个消费者协作领取同一 topic 的记录,消费者数可以超过 partition 数,并能逐条确认、释放或拒绝;传统 Consumer Group 更适合保留 partition 内顺序和 offset 语义。我要为每条任务设计幂等键,按处理时延设置锁和重试策略,监控获取、确认、释放、拒绝和超时,并为毒性消息设置隔离路径。若业务依赖 key 顺序、事务边界或客户端尚未支持,我保留 Consumer Group,先做小流量双读或独立 topic 验证,再决定迁移。”

分步骤深入解答

1. 先画出两种分配模型

Consumer Group 通常把 partition 分配给成员;同一 partition 在一个 group 内由一个成员读取,因此并行度受 partition 数约束。Share Group 则让成员从订阅的 topic 中协作领取记录,一个 partition 可以同时有多个成员处理不同记录,成员数可以超过 partition 数。这个差异适合独立工作项,却不能直接推导出全局顺序。

text
Consumer Group:  partition-0 -> worker-A
                 partition-1 -> worker-B
                 extra workers wait for another partition

Share Group:     partition-0 records -> worker-A, worker-B, worker-C
                 each acquired record is locked for one consumer

选择 Share Group 的理由应是“需要弹性领取和逐条完成”,而不是“partition 太少所以一定要换”。如果同一 customer 的事件必须依次生效,可以继续按 key 分区并使用 Consumer Group,或在业务层建立串行化队列。

2. 把记录生命周期写成状态机

KIP-932 描述了带时间限制的 acquisition lock。消费者取得记录后可以 acknowledge 表示成功、release 让记录再次可投递、reject 表示不可处理,或者什么都不做等待锁到期。默认锁时长在 KIP 中为 30 秒,但部署应以实际 broker 配置为准,不能把这个默认值当成 SLA。

text
available -> acquired -> acknowledged
                    -> released -> available
                    -> rejected  -> terminal or quarantine
                    -> lock timeout -> available

处理函数必须先用业务幂等键登记意图,再执行外部副作用;否则锁超时或客户端崩溃会造成重复扣款、重复发货或重复通知。确认成功只能代表这次领取完成,不能替业务系统回滚已经提交的副作用。

3. 设计重试、毒性消息和并发上限

delivery attempt 计数可用于区分短暂故障和不可处理记录。对网络抖动可 release 并使用退避;对 schema 不兼容、必填字段缺失等确定性错误,应 reject 并写入隔离 topic 或人工队列。不要无限 release,否则单条毒性消息会持续消耗处理预算。

锁时长应覆盖正常处理 p99 加上可解释的抖动余量;过短会导致并发重复获取,过长会拖慢恢复。还要限制每个 partition 的 acquired record 数,配合 worker semaphore、数据库连接池和外部 API 配额。监控必须同时展示 active locks、lock timeout、attempt 分布、reject 数和端到端完成延迟。

4. 重新定义顺序和重复语义

传统 Kafka 叙述经常把“partition 内有序”误当成“业务处理有序”。Share Group 允许多个消费者并行领取,同一 key 的完成顺序可能与写入顺序不同;释放、重新投递和不同批次处理还会放大这种差异。若业务需要顺序,必须把 key 级串行化、版本号检查或状态机约束写进应用,而不是只在面试中说“Kafka 有序”。

Exactly-once 也不能凭 group 类型自动获得。需要把消息读取、业务写入和确认边界逐项核验;外部数据库或支付系统仍要依赖幂等键、去重表或事务型 outbox。对不支持的组合,明确采用 at-least-once 加幂等,而不是声称“Share Group 就是 exactly-once”。

5. 规划迁移和回滚

先确认 broker 版本、客户端 API、group 类型配置、ACL、监控和运维命令,再在独立 topic 或小规模工作负载上压测。用故障注入验证:worker 在获取后崩溃、处理超过锁时长、连续 reject、broker 重启和 coordinator 切换。记录每条记录的业务键、attempt、状态和时间线。

若旧消费者仍依赖顺序或事务,不能直接把同一个 group 原地改成 Share Group。更安全的方式是复制到专用工作 topic,让新 group 逐步接管;保持旧 group 可回放,达到错误率、重复副作用和延迟门槛后再扩大流量。回滚应停止新 group 的领取并让旧路径继续消费尚未迁移的记录,避免两条路径同时执行同一副作用。

高质量示范回答

“我会先问工作项能否乱序、处理是否幂等,以及现有部署是否支持 KIP-932。Share Group 适合把 topic 中的独立记录当作协作工作项:多个成员可以从同一 partition 领取不同记录,成员数不再受 partition 数直接限制,并且每条记录有 acknowledge、release、reject 和锁超时路径。它牺牲或改变了传统 group 的分配与顺序直觉,所以同一业务 key 若要求顺序,我会继续使用 Consumer Group 或增加应用层版本检查。

我会为每条记录设置业务幂等键,按 p99 处理时间配置 acquisition lock,限制 active locks 和外部依赖并发;短暂错误退避 release,确定性坏数据 reject 到隔离路径,超过 attempt 阈值停止自动重试。监控获取、确认、释放、拒绝、锁超时、重复副作用和完成延迟。迁移前做版本、ACL、客户端和故障注入验证,先用独立 topic 小流量运行。除非证明消息读取、业务写入和确认共享同一事务边界,否则我会明确采用 at-least-once 加幂等,不把 Share Group 宣称为 exactly-once。”

常见错误

  • 把 Share Group 说成 RabbitMQ 克隆 → 两者的存储、回放和管理模型不同 → 只承诺 KIP 明确的协作领取与确认语义。
  • 按 partition 数设置最大 worker 数 → Share Group 的并行模型允许多个成员处理同一 partition → 按锁、下游容量和端到端延迟限制并发。
  • 默认保留 key 顺序 → 多成员领取和重试会改变完成顺序 → 对需要顺序的 key 使用串行化或版本检查。
  • 处理成功后直接确认但没有幂等 → 崩溃或锁超时会再次投递 → 先以业务键去重,再确认记录。
  • 无限 release 毒性消息 → 重试会占满锁和下游预算 → 按错误类型、attempt 和隔离策略终止自动重试。
  • 把默认 30 秒当作保证 → broker 配置和处理延迟可能不同 → 按实际配置与 p99 测试锁时长。

追问及应对

如果同一订单的事件必须严格按顺序怎么办?

不要直接采用 Share Group。保留以订单 ID 分区的 Consumer Group,或让同一订单进入应用层串行状态机;若必须共享领取,则需要版本号、前置版本检查和失败重排,并承认这增加了复杂度。

一个消费者拿到记录后卡住 2 分钟怎么办?

设置略高于正常 p99 的锁时长,监控 lock timeout;超时后允许再次投递,但业务处理必须幂等。对长任务可拆成可恢复步骤或外置 lease,而不是无限延长锁并隐藏故障。

如何处理连续五次 schema 错误?

把它视为确定性错误,达到阈值后 reject 并写入隔离 topic,保留原始 payload、schema 版本和错误原因。修复消费者后再通过受控重放恢复,不能让主工作流无限重试。

现有系统依赖 Kafka 事务,能否直接切换?

先列出事务覆盖的读取、处理和写入边界,再验证 Share Group 客户端与事务 API 的实际支持。若外部副作用不在同一事务内,就采用 outbox、幂等键和补偿流程;没有证据时保留 Consumer Group。

如何证明迁移没有制造重复扣款?

用唯一业务键和去重约束记录每次执行,注入 worker 崩溃、锁超时、重试和回滚故障,比较执行次数与确认次数。只有副作用次数符合业务不变量、延迟和错误率门槛,才扩大流量。

公开来源

同类题目