Kafka Consumer 深入 📖 前置阅读:本文假设读者已掌握 SpringBoot Kafka 的基本消费操作(@KafkaListener)。如果还不熟悉,建议先阅读 SpringBoot Kafka 全操作指南。
一、⚡ 问题切入:消费者重启后,怎么知道上次读到哪了? RabbitMQ 的答案是"消息消费后就删了,不需要记位置"。RocketMQ 的答案是"Broker 帮你记 offset"。Kafka 的答案是——消费者自己记,记在一个叫 __consumer_offsets 的内部 Topic 里:
消费者在 Partition-2 上消费到 offset=1500 ↓ 提交 offset __consumer_offsets Topic: Key: (order-consumer-group, order-topic, 2) Value: offset=1500 消费者重启 ↓ ↓ 读取 offset 从 offset=1501 继续消费 这个设计是 Kafka 和 RabbitMQ/RocketMQ 最核心的消费端差异——Kafka 的消费者对自己的消费进度负全责。如果消费者忘记提交 offset,重启后就会从上次提交的位置重新消费,产生重复消息。
二、Offset 提交机制 2.1 自动提交 vs 手动提交 Kafka 提供了两种 Offset 提交方式:
提交方式 配置 行为 风险 自动提交 enable-auto-commit: true 每隔 auto.commit.interval.ms(默认 5s)自动提交 poll 返回的最大 offset 消息可能没处理完就提交了——进程挂了会丢消息 手动提交 enable-auto-commit: false + ack-mode: manual 消费者处理完消息后显式调用 ack.acknowledge() 消息可能处理完了但没提交——重启后重复消费 flowchart TD classDef startEnd fill:#701a4c,stroke:#e11d48,stroke-width:2px,color:#fce7f3,font-weight:bold; classDef condition fill:#2a1147,stroke:#a855f7,stroke-width:1.5px,color:#ede9fe,font-weight:bold; classDef process fill:#1e1e24,stroke:#6b7280,stroke-width:1.5px,color:#e5e7eb; classDef highlight fill:#450a0a,stroke:#dc2626,stroke-width:1.5px,color:#fecaca,font-weight:bold; POLL([poll 拉取消息]) --> PROCESS[处理消息] PROCESS --> MODE{提交模式} MODE -- "自动提交" --> AUTO["每隔 auto.commit.interval.ms\n自动提交最后一次 poll 的 offset"] MODE -- "手动提交" --> MANUAL["业务处理成功后\n显式调用 ack.acknowledge()"] AUTO --> RISK1["风险:消息还没处理完\n但 offset 已提交\n→ 进程挂了丢消息"] MANUAL --> RISK2["风险:消息已处理完\n但 offset 没提交\n→ 重启后重复消费"] class POLL startEnd; class MODE condition; class AUTO,MANUAL highlight; class RISK1,RISK2 process; 手动提交比自动提交更安全——至少你知道什么时候提交了。消息重复消费可以用幂等解决,但消息丢失无法恢复。
...