Skip to content

Apache Kafka 可靠性

本页结论:Kafka 的可靠性由三处配置共同决定——生产端 acks + 幂等/事务、Broker 端副本与 ISR、消费端 offset 提交时机。「提交 offset」与「业务提交」是两个独立动作,之间的崩溃窗口靠幂等消费兜底;事务只覆盖 Kafka 内部,不等于外部系统 exactly-once。

生产端:acks、重试与幂等

acks确认条件丢失窗口
0不等确认网络/Broker 故障即丢
1Leader 写入内存/页缓存Leader 崩溃且未同步到 Follower → 丢
all(-1)ISR 全部副本写入取决于 ISR 收缩与 min.insync.replicas

配套关系(规格 §7.2 强制点):

  • acks=all 必须配合副本数 ≥2 才有意义;ISR 只剩 Leader 时 all 退化为 1。min.insync.replicas 让 Broker 在 ISR 不足时拒绝写入(宁可不可用,不假确认)。
  • Producer 重试:网络抖动时客户端自动重发,非幂等配置下可能产生重复记录(同内容两个 offset)。
  • Idempotent Producerenable.idempotence=true,本仓库默认):为每条记录带 producer id + 序列号,Broker 去重,重试不再产生重复。它是事务的前置条件。
java
props.put("acks", "all");
props.put("enable.idempotence", "true");
RecordMetadata md = producer.send(record).get(); // partition + offset

消费端:offset 提交即「确认」

Kafka 没有单条 ACK;消费进度就是已提交 offset。

模式行为风险
自动提交(enable.auto.commit=true)按时间间隔提交「已 poll 到」的位点处理中崩溃 → 已提交但未处理完 = 丢;提交后重新 poll 前崩溃边界附近 = 重复
手动提交(commitSync/commitAsync)业务决定提交时机提交早于业务完成同样有丢窗口

崩溃窗口与幂等消费(§5.4 基准实现)

正确顺序:业务事务提交 → commitSync

text
1. 开启本地数据库事务
2. 插入 messageId 到 processed_messages(唯一键)
3. 唯一键冲突 → duplicate_skipped
4. 首次处理 → 执行业务写入,提交事务
5. 提交成功后才 commitSync(含该分区新位点)

第 4 步成功、第 5 步前崩溃 → 重启后从旧 offset 重读 → 幂等表拦截。该窗口不可消除。「提交 offset 等于业务数据库已成功提交」是规格 §7.2 的禁止表述——两者是独立系统上的两个动作。本仓库 kafka basic 实验即按此实现:

verifiedkafka / basicbroker 4.3.1 · java-kafka-clients-4.3.1
镜像apache/kafka:4.3.1@sha256:77e3df9054047a88b520d0cc46e16696d3b22022e1d580aeccd2632df6532837
捕获时间2026-08-19T09:27:00.562Z
耗时 / 退出码27918 ms / exit 0
断言
produced3
received3
uniqueMessageIds3
businessCommitted3
business_rows3
consumerGroupLag0
consumerExitCode0
归一化日志
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.basic partitionOrQueue=0 offset=0 seq=1 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.basic partitionOrQueue=0 offset=1 seq=2 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.basic partitionOrQueue=1 offset=0 seq=3 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic destination=orders.basic produced=3 status=done
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic destination=orders.basic consumerGroup=orders-basic-group consumer=consumer-1 partitions=0,1,2 status=assigned
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.basic partitionOrQueue=0 offset=0 seq=1 consumerGroup=orders-basic-group consumer=consumer-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.basic partitionOrQueue=0 seq=1 consumerGroup=orders-basic-group attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.basic partitionOrQueue=0 offset=1 seq=2 consumerGroup=orders-basic-group consumer=consumer-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.basic partitionOrQueue=0 seq=2 consumerGroup=orders-basic-group attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.basic partitionOrQueue=1 offset=0 seq=3 consumerGroup=orders-basic-group consumer=consumer-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.basic partitionOrQueue=1 seq=3 consumerGroup=orders-basic-group attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic destination=orders.basic consumer=consumer-1 received=3 status=done
[inspect] timestamp=<ts> level=INFO service=order-service product=kafka lab=basic business_rows=3 processed_rows=3 status=snapshot
[assert] produced=3 PASS
[assert] received=3 PASS
[assert] uniqueMessageIds=3 PASS
[assert] businessCommitted=3 PASS
[assert] business_rows=3 PASS
[assert] consumerGroupLag=0 PASS
[assert] consumerExitCode=0 PASS
npm run lab -- kafka basic

事务的边界

Transactional Producer 把一个事务内的多条 produce(可加 sendOffsetsToTransaction 把消费位点一起原子提交)变成原子可见:

  • read_committed 消费者只能看到已 commitTransaction 的消息;abortTransaction 的消息永不可见。
  • 幂等 producer 保证序列号无重复;事务协调器(transaction coordinator)管理事务状态。

动手验证(3 条 commit 事务可见,2 条 abort 事务不可见):

bash
npm run lab -- kafka idempotence-transaction
verifiedkafka / idempotence-transactionbroker 4.3.1 · java-kafka-clients-4.3.1
镜像apache/kafka:4.3.1@sha256:77e3df9054047a88b520d0cc46e16696d3b22022e1d580aeccd2632df6532837
捕获时间2026-08-19T09:30:45.243Z
耗时 / 退出码29236 ms / exit 0
断言
txnCommitted1
txnAborted1
committedVisible3
abortedVisible0
business_rows3
consumerGroupLag0
consumerExitCode0
归一化日志
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.txn partitionOrQueue=0 offset=0 seq=1 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.txn partitionOrQueue=0 offset=1 seq=2 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.txn partitionOrQueue=1 offset=0 seq=3 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction destination=orders.txn produced=3 status=txn_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction destination=orders.txn consumerGroup=orders-txn-group consumer=consumer-1 partitions=0,1,2 status=assigned
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.txn partitionOrQueue=0 offset=0 seq=1 consumerGroup=orders-txn-group consumer=consumer-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.txn partitionOrQueue=0 seq=1 consumerGroup=orders-txn-group attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.txn partitionOrQueue=0 offset=1 seq=2 consumerGroup=orders-txn-group consumer=consumer-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.txn partitionOrQueue=0 seq=2 consumerGroup=orders-txn-group attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.txn partitionOrQueue=1 offset=0 seq=3 consumerGroup=orders-txn-group consumer=consumer-1 attempt=1 redelivered=false status=received
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-3 eventType=order.created schemaVersion=1 aggregateId=order-1003 traceId=trace-3 correlationId=order-1003 destination=orders.txn partitionOrQueue=1 seq=3 consumerGroup=orders-txn-group attempt=1 status=business_committed
[consumer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction destination=orders.txn consumer=consumer-1 received=3 status=done
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-1 eventType=order.created schemaVersion=1 aggregateId=order-1001 traceId=trace-1 correlationId=order-1001 destination=orders.txn partitionOrQueue=0 offset=3 seq=1 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction messageId=mid-2 eventType=order.created schemaVersion=1 aggregateId=order-1002 traceId=trace-2 correlationId=order-1002 destination=orders.txn partitionOrQueue=0 offset=4 seq=2 durationMs=<ms> status=produced
[producer] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction destination=orders.txn produced=2 status=txn_aborted
[inspect] timestamp=<ts> level=INFO service=order-service product=kafka lab=idempotence-transaction business_rows=3 processed_rows=3 status=snapshot
[assert] txnCommitted=1 PASS
[assert] txnAborted=1 PASS
[assert] committedVisible=3 PASS
[assert] abortedVisible=0 PASS
[assert] business_rows=3 PASS
[assert] consumerGroupLag=0 PASS
[assert] consumerExitCode=0 PASS
npm run lab -- kafka idempotence-transaction

边界必须说清楚(禁止表述之三):

  • Kafka EOS 的 exactly-once 指 Kafka 内部(topic → topic 的 produce-consume)不重不丢;
  • 开启事务后,任意外部系统副作用(数据库、HTTP、短信)都不是 exactly-once——外部写入仍需幂等设计。

顺序与重试的关系

  • 分区内顺序由写入顺序决定;Producer 开 max.in.flight.requests.per.connection>1 且非幂等时,重试可能乱序。幂等 producer 允许最多 5 个在途请求仍保序。
  • 消费端「失败重试」通常意味着跳过或转发该消息(应用层 DLQ),没有 RabbitMQ 式 requeue 到队头。

三层语义总结

层级保证条件
Broker确认的消息在 ISR 全部副本可读acks=all + min.insync.replicas≥2(多副本);本仓库单节点 RF=1 仅覆盖 Leader 存活场景
Client不重发丢、不提前提交幂等 producer + 异常重试;手动 commitSync 且晚于业务提交
Business效果恰好一次幂等表 + 本地事务;外部系统副作用不在 Kafka 事务覆盖内

常见误区

  • 「acks=all 就绝对不丢」——ISR 收缩到 1 时确认不再要求第二个副本;必须与 min.insync.replicas 一起配置。
  • 「事务提交后所有下游都恰好一次」——见上文边界说明。
  • 「消费者崩溃只重投一条」——重投单位是「已提交 offset 之后的全部消息」。

官方资料

以统一实验验证消息系统语义边界