消费一条消息后还要产生新消息,如何避免处理中断造成重复或丢失?
我会先控制影响,再按证据定位:幂等生产者通过 Producer ID、序列号和 Broker 去重,避免单会话重试造成分区内重复;Kafka 事务可把多分区写入和消费位点提交放进同一事务。消费者使用 readcommitted 时只读取已提交事务,从而支持 Kafka 内部处理链路的 Exactly-Once。
我会先确认影响范围,同时控制故障继续放大。理解幂等生产者、事务写入、隔离级别及端到端恰好一次的边界。幂等生产者通过 Producer ID、序列号和 Broker 去重,避免单会话重试造成分区内重复;Kafka 事务可把多分区写入和消费位点提交放进同一事务。消费者使用 readcommitted 时只读取已提交事务,从而支持 Kafka 内部处理链路的 Exactly-Once。 Exactly-Once 是有范围的语义,不代表发送邮件、调用支付接口或写普通数据库也会自动恰好执行一次。
止损之后,我会按请求链路建立证据,而不是凭经验猜。若事务提交失败或生产者崩溃,输出记录对 readcommitted 消费者不可见,来源位点也不会推进,重启后可安全重新处理。它将“输出写入”和“输入进度”绑定在 Kafka 内部,而不是把任意外部副作用放进一个魔法全局事务。 TransactionCoordinator:事务状态与 marker。 ProducerStateManager:PID/epoch fencing。 我会先确认请求实际走到了哪条路径,再用运行数据验证,不会只看类名或配置猜测。
定位时我最关注这些参数、指标和容量关系。在 sendOffsetsToTransaction 前后 kill 进程,核对输入位点与输出原子性。
找到根因后先做最小修复,再用同样的流量验证。Kafka Streams 输出 Topic 使用 EOS,随后消费者调用数据库普通 INSERT,重放时仍重复。EOS 边界只覆盖 Kafka 内部;数据库需唯一业务键、幂等 upsert 或 Outbox/CDC 协调。 把 Kafka EOS 延伸解释到数据库:事务提交/中止率:稳定唯一 transactional.id。 固定输入和基线 先在没有故障注入的环境执行上述配置,固定数据规模、并发度、运行时版本和预热时间。以「transaction abort」为主基线,记录值应满足「记录分区基线」;同时保存 事务提交/中止率、producer fencing,使后续变化能够回到同一时间轴比较。
恢复阶段要逐步放量,最后把监控和边界补齐。方案:更适合的场景:主要收益:代价与边界。 幂等生产者:只防生产重试写重复:默认易用、开销低:不原子绑定多个分区与消费位点。 Kafka 事务/EOS:consume-transform-produce 全在 Kafka:位点与输出原子:事务开销、超时与隔离级别复杂。 业务幂等/Outbox:副作用跨数据库或外部系统:覆盖最终资源:需要状态表、重试和对账。