Spring Kafka 与消息集成

共 19 题
📑 题目列表 19 题
#
★★★

1. Kafka 事务能原子提交多个分区记录和消费偏移,为何不能直接保证数据库更新也原子

Kafka 事务能原子提交多个分区记录和消费偏移(跨分区/主题),但为何不能直接保证数据库更新也原子?

  • Kafka 事务的原子范围(Kafka 内部)
  • 数据库与 Kafka 是不同系统
  • 跨系统事务的局限(Outbox)

Kafka 事务(enable.idempotence + transactional.id)保证的是"Kafka 内部的原子性":消息写入多个分区/主题、以及消费偏移的提交,在同一个事务内原子可见(read_committed 消费者只读已提交消息)。但它无法保证"数据库更新"与"Kafka 消息发送"作为整体原子,因为数据库和 Kafka 是两套独立的存储系统,没有共享的事务协调器(XA/两阶段提交在跨异构系统上复杂且不常用)。Kafka 事务只协调 Kafka 内部(broker + 事务协调器),无法介入数据库的提交。因此"业务更新 DB + 发消息"要原子,需用 Outbox Pattern(本地消息表:DB 事务内写消息表,Canal/定时任务投递到 Kafka)。

事务原子性要"同一事务管理器"范围内。Kafka 事务是 Kafka 内部事务,与数据库事务是两个独立事务,无法跨系统原子。Outbox 用"本地 DB 事务写消息表 + 异步投递"实现最终一致。

#
★★★

2. Kafka 的 JMX 指标与 Prometheus 集成

Kafka 的 JMX 指标如何与 Prometheus 集成?

  • Kafka 暴露 JMX 指标
  • JMX Exporter 转换
  • Prometheus 抓取与告警

Kafka(broker 与客户端)通过 JMX 暴露指标(如 kafka.server:type=BrokerTopicMetrics 的每秒消息数、UnderReplicatedPartitions、请求吞吐、#Partitions 等)。与 Prometheus 集成:用 JMX Exporter(jmx_prometheus_javaagent)作为 Java agent 挂载到 Kafka,把 JMX 指标暴露为 Prometheus 的 HTTP 抓取端点,Prometheus 定时抓取并存储,配合 Grafana 面板与告警规则。关键指标:UnderReplicatedPartitions(副本不足)、RequestsPerSecBytesIn/OutPerSecISR 状态、消费者 group 的 lag。集成后实现 Broker 与消费者监控告警。

集成链路是"JMX 指标 -> JMX Exporter -> Prometheus -> Grafana"。监控重点是副本健康、吞吐、消费积压与延迟。

#
★★★

3. Kafka 的 Schema Registry 与 Avro/JSON Schema

Kafka 的 Schema Registry 与 Avro/JSON Schema 如何配合?

  • Schema Registry 管理消息 schema
  • Avro/JSON Schema 序列化
  • 兼容性演进

Schema Registry 集中管理 Kafka 消息的 schema(Avro/JSON Schema/Protobuf),生产者在发送前向 Schema Registry 注册 schema 并获取 schema id,消息头携带 schema id;消费者根据 id 从 Registry 取 schema 反序列化。结合 Avro 等二进制格式,实现高效序列化与强类型。核心价值:schema 版本管理与兼容性演进——Registry 支持 BACKWARD/FORWARD/FULL 兼容性检查,阻止不兼容变更,避免生产消费者因 schema 变化解析失败。Spring Kafka 通过 KafkaAvroSerializerio.confluent)或 SchemaRegistry 组件集成。JSON Schema 也可用,但 Avro 更紧凑。

Schema Registry 解决"消息 schema 谁来管、如何演进"的问题,Avro 提供紧凑二进制与 schema 演进,Registry 保证兼容性。

#
★★★

4. Kafka 的 Topic/Partition/Offset 模型

Kafka 的 Topic/Partition/Offset 模型是怎样的?

  • Topic 逻辑队列
  • Partition 分区与顺序
  • Offset 位移

Kafka 以 Topic 组织消息,每个 Topic 分为多个 Partition(物理分区),消息按 offset(从 0 递增的位移)追加到分区日志。分区内有序(按 offset),但跨分区不保证全局顺序。Partition 是并行与复制的最小单位:每个分区有独立副本(leader/follower),消费者组内每个分区同一时刻只被一个消费者消费(保证分区内有序)。Offset 是消息在分区内的位置,消费者提交 offset 记录消费进度,重放可 seek 到指定 offset。Topic 的并行度 = 分区数,吞吐随分区扩展。

分区模型是 Kafka 的基石:分区内有序 + 分区可并行 + 副本再分区,offset 支持消费进度与重放。

#
★★★

5. Kafka 的 enable.auto.commit 与手动提交

Kafka 的 enable.auto.commit 与手动提交有何区别?

  • 自动提交与手动提交
  • 提交时机与语义
  • ack 与可靠性

enable.auto.commit=true(默认)时,消费者定时自动提交 offset(auto.commit.interval.ms),实现简单但可能丢消息(提交后处理失败)或重复(处理中崩溃未提交,重新消费)。enable.auto.commit=false 则手动提交:commitSync() 同步阻塞提交(确认成功)、commitAsync() 异步提交(不阻塞但可能失败,需回调处理)。手动提交让开发者控制"处理成功后再提交",保证 at-least-once(处理完才提交)。Spring Kafka 中 AckModeenable.auto.commit 配合,MANUAL 模式需手动 ack。可靠性:手动提交(处理成功后再提交)更可靠,但需处理重复(幂等)。

自动提交"拿来即提交",手动提交"处理完再提交"。手动提交换来 at-least-once 与业务控制,代价是需处理重复与提交失败。

#
★★★

6. Kafka 的 isolation.level(read_uncommitted/read_committed)

Kafka 的 isolation.level(read_uncommitted/read_committed)有何区别?

  • 隔离级别语义
  • 事务消息可见性
  • 与 Kafka 事务配合

isolation.level 决定消费者如何读取事务消息:read_uncommitted(默认)读取所有消息,包括未提交的(可能被事务回滚,读到"幻影"消息);read_committed 只读取已提交事务的消息,且能读到事务的"结束标记"(abort/txn marker),跳过被回滚的消息。使用 read_committed 需要消费者配合 Kafka 事务(生产者 enable.idempotence + 事务),才能获得"只读已提交"的语义。read_committed 会带来轻微延迟(等待事务结束标记)与更高的顺序开销。选型:普通非事务场景用默认 read_uncommitted;需要事务消息原子可见性用 read_committed。

isolation.level 控制事务消息的可见性边界,read_committed 只读已提交,避免读到回滚数据,适合事务消息场景。

#
★★★

7. Kafka 的 max.poll.interval.ms 与 session.timeout.ms

Kafka 的 max.poll.interval.ms 与 session.timeout.ms 有何区别与作用?

  • max.poll.interval.ms 处理超时
  • session.timeout.ms 心跳超时
  • rebalance 触发

max.poll.interval.ms 是消费者两次 poll() 之间的最大间隔,超过则消费者被认为"处理太慢",主动触发 rebalance(把分区转给其他消费者)。session.timeout.ms 是消费者与 broker 的心跳超时,超过该时间未收到心跳,broker 判定消费者死亡,触发 rebalance。区别:max.poll.interval.ms 面向"处理慢/卡死"(业务处理超时),session.timeout.ms 面向"心跳中断/网络故障"(消费者失联)。两者都触发 rebalance。处理慢可调大 max.poll.interval.ms 或减少单次 poll 消息数(max.poll.records)、异步化处理。session.timeout.ms 需配合 heartbeat.interval.ms

两个超时一个管"处理慢"一个管"心跳失联",都触发 rebalance。处理慢要调大 max.poll.interval.ms 或拆分处理,避免频繁 rebalance。

#
★★★

8. Kafka 消费积压时增加消费者的上限为什么是分区数,扩分区对已有消息顺序的影响

Kafka 消费积压时增加消费者的上限为什么是分区数,扩分区对已有消息顺序有何影响?

  • 消费者数 <= 分区数
  • 分区内有序
  • 扩分区与顺序

消费者组内,一个分区同一时刻只能被组内一个消费者消费(保证分区内有序),因此消费者数量上限 = 分区数;消费者数超过分区数时多出的消费者空闲。消费积压时,若消费者数已达到分区数,增加消费者无法提升吞吐,需增加分区数(扩分区)才能提高并行度。扩分区对已有消息顺序的影响:新消息进入新分区,但已有分区内消息顺序不变;跨分区本来就不保证全局顺序,因此扩分区不破坏"分区内有序"保证,但会改变"原打算按写入顺序的全局顺序"(Kafka 本不含全局顺序)。扩分区后需注意 hash 分区模式下 key 到分区的映射变化,可能影响按 key 路由。

并行度上限 = 分区数,消费者数 <= 分区数。扩分区是提升并行度的手段,分区内有序保持,全局顺序本就不保证。

#
★★

9. Spring Kafka 4.0 中 KafkaTemplate.send 与 CompletableFuture 异步回调在虚拟线程下的阻塞语义

Spring Kafka 4.0 中 KafkaTemplate.send 与 CompletableFuture 异步回调在虚拟线程下的阻塞语义如何?

  • KafkaTemplate.send 返回 CompletableFuture
  • 异步回调与虚拟线程
  • 阻塞 vs 异步

KafkaTemplate.send(...) 返回 CompletableFuture<SendResult>,底层 KafkaProducer.send 是异步的(消息入队后返回,回调在单独线程触发)。在虚拟线程下:send 本身不阻塞(立即返回 Future),调用方虚拟线程可继续执行;若在虚拟线程上调用 future.get() 阻塞等待结果,虚拟线程会挂起不占平台线程,符合虚拟线程模型。回调(whenComplete/thenApply)在 Kafka 的 producer 回调线程执行,不阻塞虚拟线程。注意:批次满时 send 可能阻塞(buffer.memory 满),虚拟线程挂起等待;需合理配置 buffer.memorymax.block.ms 避免虚拟线程堆积。异步回调要避免在回调里做阻塞操作。

send 是异步非阻塞,虚拟线程下阻塞等待 (get) 挂起虚拟线程不占平台线程,但 buffer 满时会阻塞,需配置缓冲与超时。

#
★★

10. Spring Kafka 的 @KafkaListener 配合 BatchListener 在 Spring Boot 4.0 虚拟线程下的执行模型

Spring Kafka 的 @KafkaListener 配合 BatchListener 在 Spring Boot 4.0 虚拟线程下的执行模型如何?

  • @KafkaListener 消费模型
  • BatchListener 批量消费
  • 虚拟线程执行

@KafkaListener 默认单条消费(MessageListener),BatchListener=true 时 batch 消费(一次 poll 返回多条 List<ConsumerRecord>)。执行模型:Spring Kafka 的 KafkaMessageListenerContainerConsumerRecordListener 分发,设置在虚拟线程(Spring Boot 4.0 支持虚拟线程执行器)时,每个 batch 在虚拟线程上处理,虚拟线程挂起不占平台线程,提升并发吞吐。注意:批量消费需 batchListenerconcurrency 配置;虚拟线程下并发 batch 处理共享同一 consumer,需注意线程安全与 offset 管理(batch 处理完统一提交)。批量 + 虚拟线程可提升吞吐,但需处理消息顺序(同一 partition 内)与提交语义。

BatchListener 批量拉取 + 虚拟线程并发处理,提升吞吐;但共享 consumer 的 offset 提交与顺序需谨慎,虚拟线程不改变 Kafka 分区内有序语义。

#
★★

11. Kafka 的分区再均衡(Rebalance)协议(Eager/Cooperative)

Kafka 的分区再均衡(Rebalance)协议(Eager/Cooperative)有何区别?

  • Eager 全量停止再分配
  • Cooperative 增量协调
  • 用 Static Group 与无状态

Eager Rebalance(旧协议):消费者进入 rebalance 时全部撤销分区(Stop-the-world),永久性中断消费,再重新分配,代价是"全体暂停"。Cooperative Rebalance(增量协调,KIP-429,Kafka 2.4+):消费者先释放部分分区(revoke),再分配新分区,可分多次收敛,避免全体暂停,减少"stop-the-world"窗口,提升稳定性。Cooperative 聚合适用于 CooperativeStickyAssignor。选用 Cooperative 能减少 rebalance 造成的消费中断,但需消费端幂等(分区可能短暂转移)。工程上优先 Cooperative,配合 session.timeoutmax.poll.interval 调优。

Eager 全体暂停再分配,Cooperative 增量协调减少暂停窗口。Cooperative 更优但需处理分区短暂转移的幂等。

#
★★

12. Kafka 的幂等 Producer(enable.idempotence=true)

Kafka 的幂等 Producer(enable.idempotence=true)如何工作?

  • 幂等 Producer 的机制
  • 去重与重试
  • 单分区幂等保证

enable.idempotence=true 启用幂等 Producer,通过 producerId + sequence(每个分区递增序号)实现:broker 校验序号,重复的序号(重试导致的重复发送)被拒绝,从而保证"单分区内、单会话内"不重复写入。它配合 acks=all(默认)与 retries 提供 exactly-once 写入(至少一次的去重)。幂等只保证单分区内不重复,跨分区/跨主题需事务(transactional.id)。启用幂等后 producer 自动处理重试导致的重复,减少"至少一次"语义下的重复。注意:producer 重启后 producerId 变化,幂等性仅限同一 producerId 会话。

幂等 Producer 用 PID+sequence 去重,解决重试重复,保证单分区不重复。跨分区/事务需结合事务 API。

#
★★

13. Kafka 的消息顺序性与单分区保证

Kafka 的消息顺序性如何保证,单分区保证是什么?

  • 分区内有序
  • 消息顺序与 key
  • 消费者并行与顺序

Kafka 只保证"分区内有序":同一分区的消息按写入顺序(offset 递增)被消费,生产者按 key 哈希路由到分区,同一 key 的消息进入同一分区,从而保证"同一 key 相关消息有序"。要保证某实体消息有序,需用该实体 id 作为 key(hash 到同一分区)。跨分区不保证全局顺序。消费端:消费者组内同一分区被一个消费者串行消费(分区内有序),若消费者并行消费同一分区(多线程)需自行保序。因此高 QPS 且保序场景,用 key 路由到分区 + 分区内单消费者。

Kafka 顺序粒度是"分区",通过 key 路由让相关消息进同一分区保序。跨分区/多线程消费需自行处理顺序。

#
★★

14. 消费者 seek/assign 手动管理位移的场景(重放、死信重试)与风险

消费者 seek/assign 手动管理位移的场景(重放、死信重试)与风险是什么?

  • seek 到指定 offset 重放
  • assign 手动分配分区
  • 风险与位移管理

seek(partition, offset) 让消费者从指定 offset 重新消费(重放),用于:修复脏数据后重放、死信重试(重新消费失败消息)、时间窗口回放(seekToBeginning/seekToTimestamp)。assign(partition) 手动指定消费的分区(不经 rebalance),用于精确控制分区。风险:手动管理位移易出错——seek 到错误 offset 导致重复或丢失;assign 后消费者不参与 rebalance,需自行管理分区与位移;手动提交与 seek 的组合需谨慎,避免跳过未处理消息。Spring Kafka 可用 SeekUtils/ConsumerSeekAware 实现 seek。核心是"位移是可靠性的锚点",手动操作需确保幂等。

seek/assign 提供灵活的重放与分区控制,但把位移管理交给开发者,风险高。需配合幂等消费与精确位移控制。

#
★★

15. Spring Kafka 的 @RetryableTopic 与死信处理器(DltHandler)如何实现重试与告警

Spring Kafka 的 @RetryableTopic 与死信处理器(DltHandler)如何实现重试与告警?

  • @RetryableTopic 重试主题
  • DltHandler 死信处理
  • 重试/告警机制

@RetryableTopic 注解让消费失败的消息自动进入重试主题(如 topic-retry-0topic-retry-1,按 backoff 递增),重试次数用尽后进入死信主题(topic-dlt)。@DltHandler 标注死信处理方法,处理无法消费的消息(记录、告警、人工处理)。使用 @RetryableTopic 需配置 backoff(间隔)、attempts(次数)、autoStartDltHandler。机制:重试主题按次数分多个,失败消息按次数路由到对应重试主题并延迟消费,最终失败进 DLT。告警:在 DltHandler 中记录日志/发送告警,或监控 DLT 消费积压。

@RetryableTopic 提供"重试主题 + 死信主题"的自动重试与失败处理,@DltHandler 兜底记录与告警,避免无限重试阻塞。

#

16. Kafka 4.x 完全移除 ZooKeeper 后,Spring Boot 4.x + Spring Kafka 在 KRaft 模式下的适配要点

Kafka 4.x 完全移除 ZooKeeper 后,Spring Boot 4.x + Spring Kafka 在 KRaft 模式下的适配要点是什么?

  • KRaft(替代 ZooKeeper)
  • 配置与部署适配
  • client 兼容

Kafka 4.x 完全移除 ZooKeeper,元数据由 KRaft(Raft 协议)管理,broker 与 controller 合一。适配要点:集群部署不再配置 ZooKeeper,改用 KRaftcontroller 节点;客户端(Spring Kafka)层面基本无需改动(协议兼容),因客户端只连 broker,不感知 ZooKeeper/KRaft 差异。Spring Kafka 4.x 只需配置 bootstrap-servers 指向 broker 即可,无需 ZooKeeper 地址。注意:KRaft 模式下需配置 process.roles(controller/broker)、controller.quorum.voters;迁移现有 ZooKeeper 集群需工具。整体上客户端适配简单,主要是部署与运维侧变化。

KRaft 移除 ZooKeeper,客户端透明。适配重点是部署架构(controller 配置)与迁移,客户端配置不变。

#

17. Producer 的 acks 参数(0/1/-1)语义

Producer 的 acks 参数(0/1/-1)有何语义?

  • acks=0 不确认
  • acks=1 leader 确认
  • acks=-1(all)全副本确认

acks 决定 producer 何时认为发送成功:acks=0——不等待确认,立即成功,吞吐最高但可能丢消息(broker 崩溃);acks=1——leader 写入即确认,正常情况不丢,但 leader 崩溃未同步副本时可能丢;acks=-1(all)——所有 ISR 副本确认后才认为成功,最可靠不丢,但延迟高、吞吐低。选型:追求吞吐且容忍丢用 acks=0/1;追求不丢用 acks=all。配合 min.insync.replicas(最少同步副本数)与 enable.idempotence(幂等要求 acks=all)保证可靠性。acks=all 降低单点故障丢消息风险。

acks 是"可靠性与性能"的权衡。acks=all 更可靠,acks=0 最高吞吐,工程上常用 acks=all + min.insync.replicas。

#

18. Spring Kafka 的消费确认,AckMode 各取值(MANUAL/RECORD/BATCH)的语义如何?

Spring Kafka 的 AckMode 各取值(MANUAL/RECORD/BATCH)的语义是什么?

  • AckMode 各取值
  • RECORD/MANUAL/BATCH
  • 提交时机

AckMode 决定 Spring Kafka 何时提交 offset(配合 enable.auto.commit=false):RECORD——每条消息处理完立即提交(粒度最细,保证 at-least-once 但频繁提交);BATCH——每次 poll 的一批消息处理完提交(更新并发效率);MANUAL——应由业务代码手动提交(Acknowledgment.acknowledge()),最灵活;MANUAL_IMMEDIATE——acknowledge 后立即提交(不等下次 poll);TIME/COUNT——按时间/数量批量提交。AUTO(默认)在监听器返回后提交(由 poll 决定)。可靠性:RECORD/BATCH 由框架根据处理完成提交,MANUAL 由业务控制。选型:需要精确控制(如处理成功才提交)用 MANUAL。

AckMode 控制 offset 提交粒度与时机。RECORD 最细粒度、BATCH 更高效、MANUAL 由业务精确控制,选型权衡可靠性与提交开销。

#

19. Kafka 压缩算法(gzip/lz4/zstd)的选择对吞吐与 CPU 的影响

Kafka 压缩算法(gzip/lz4/zstd)的选择对吞吐与 CPU 有何影响?

  • 各压缩算法特性
  • 压缩比 vs CPU
  • 选型

Kafka 支持 gzipsnappylz4zstd 压缩。lz4/snappy 压缩速度快、CPU 开销低、压缩比适中,适合追求吞吐的低 CPU 场景;gzip 压缩比高但 CPU 开销大、速度慢;zstd 在压缩比与速度之间平衡较好(高压缩比 + 较快速度),是现代推荐。选型权衡:数据量大、带宽是瓶颈时选高压缩比(zstd/gzip);CPU 紧张、追求低延迟时选 lz4/snappy。压缩发生在 producer 端,broker 保留压缩,consumer 端解压。需结合数据特征(文本可压缩性高)与硬件(CPU 核数)选择。

压缩是"带宽 vs CPU"的权衡。zstd 兼顾压缩比与速度,lz4 快,gzip 压缩比高但慢,按数据与硬件选型。