尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Kafka精确一次投递:幂等性、事务机制与工程实践深度解析

Kafka精确一次投递:幂等性、事务机制与工程实践深度解析 1. 为什么精确一次这么难三种投递语义的本质差异老规矩我们先把最核心的概念掰开揉碎。Kafka官方文档里定义了三种投递语义at-most-once最多一次、at-least-once至少一次、exactly-once精确一次。我当年第一次看到这三个词觉得无非是丢不丢消息、重不重复的问题但直到我在生产环境里亲手排查过一次重复订单流水之后才明白这三者背后对应的是完全不同的系统设计和性能代价。1.1 消息丢失、重复与乱序从消费者视角看问题先看一个最朴素的问题消费者到底会看到什么at-most-once模式下Producer发送消息后不等任何确认消息可能丢了就丢了消费者最多收到一次。这种模式的优点是延迟极低适用于日志采集、监控指标上报这类丢了也无所谓的场景但业务数据绝对不能用。at-least-once是Kafka默认的行为。Producer把acks设置为allBroker确认写入后才算发送成功消费者拉取消息后先处理业务、再提交offset。如果消费者在处理完消息但还没提交offset时宕机了重启后会从旧offset重新消费这就产生了重复。所以至少一次的必然代价是重复消费。exactly-once则要求每条消息对下游的影响恰好等于一次不多不少。这里要注意一个关键点Kafka本身能做到的是写入只发生一次但对于从Kafka读数据、处理、再写回Kafka的完整链路仅仅靠Broker端的机制是不够的。1.2 幂等性、事务与精确一次三者的关系和边界很多初学者把这三个概念当成三个并列的选项其实它们是层层递进的关系机制解决的核心问题作用范围幂等性IdempotenceBroker端因Producer重试导致的重复写入单个Producer会话内单个分区事务Transactions跨分区原子写入要么全部成功要么全部失败单个Producer生产的事务批次精确一次处理EOS端到端的读-处理-写链路保证消费处理生产全链路打个比方幂等性是同一笔订单我反复提交系统只认一次事务是购物车里的十件商品必须一起扣库存不能扣一半精确一次则是从下单到支付到发货整个流程在系统的每个环节都只生效一次。搞清楚了这些边界再去看Kafka的源码和配置很多疑惑都会迎刃而解。2. 幂等性机制PID与序列号的协同过滤幂等性是Kafka消息保证体系的基石但它的实现细节比大多数人以为的要巧妙得多。我见过不少同学以为设置了enable.idempotencetrue就万事大吉实际上它只在特定场景下生效踩坑之前最好先把原理吃透。2.1 幂等性到底解决了什么问题在没有幂等性之前Kafka Producer遇到网络抖动或者Broker端短暂不可用时会自动重试。重试本身是好事但会带来一个非常棘手的副作用第一次发送的消息其实已经写入成功了只是响应超时Producer不知道于是又发了一次。结果是Broker上同一个分区里出现了两条一模一样的消息。幂等性要解决的就是这个因重试导致的重复写入。它的实现逻辑不复杂Producer在初始化时从Broker申请一个全局唯一的PIDProducer ID同时为每条消息维护一个从0开始单调递增的序列号Sequence Number。Broker端收到消息时会检查(PID, 分区, 序列号)这个三元组如果序列号比之前收到的最大值大1正常写入如果序列号 已记录的最大序列号说明是重复消息直接丢弃如果序列号出现跳跃比如中间缺了几个说明有消息确实丢了Broker会返回OutOfOrderSequenceException。2.2 Broker端如何识别重复消息Broker端并不是把所有消息的序列号都存起来那样内存开销太大了。Kafka在每个分区对应的Log文件里维护了一个ProducerStateEntry只记录每个PID最近写入的5个序列号这是max.in.flight.requests.per.connection最大允许的值用于快速判断重复和乱序。有一点值得注意幂等性要求单分区内写入有序。如果max.in.flight.requests.per.connection设置大于1且没有开启幂等Broker端可能因为网络返回顺序错乱而写入乱序消息。开启幂等性之后即使in-flight请求超过1Broker端也会通过序列号机制把这些消息在分区日志中排列成正确的顺序。2.3 幂等性的局限跨分区跨会话的重复幂等性有一个天然的边界PID只在Producer进程的生命周期内有效。一旦Producer进程重启它会申请一个新的PID之前PID对应的序列号记录全部作废。这意味着同一条消息如果由旧PID和新PID各写一次Broker无法识别它们是重复的幂等性只针对同一个Producer的写入如果应用层有多个Producer实例比如分布式部署它们各自维护独立的PID跨实例的重复无法靠这个机制解决。我在一个订单系统中遇到过这样的案例服务A把订单消息发给Kafka服务B消费后做处理结果因为一次部署重启服务A的发送逻辑被重放了一遍应用层自己的重放机制不是Kafka的重试生成了两条完全相同的订单消息。幂等性在这里无能为力因为这是两条来自不同PID的消息。幂等性解决的是同一Producer重试导致的重复不是业务逻辑层面的重复。理解了这个边界你才知道什么时候该用事务、什么时候该在业务侧加幂等表。3. 事务机制跨分区原子写入的实现路径如果你需要在Kafka里实现要么全部写入成功要么全部不写入的原子性就需要用到Kafka事务。它解决的是幂等性覆盖不到的另一个问题多分区、多Topic之间的原子写入。3.1 事务协调器与事务日志Kafka的事务机制依赖一个专门的角色Transaction Coordinator事务协调器。每个Broker上都会运行一个这样的协调器组件但真正承担某个Producer事务协调工作的是那个以及分区哈希计算结果对应的Broker。事务的元数据比如事务状态、参与的分区列表会被写入一个名为__transaction_state的内部Topic。这个Topic默认有50个分区每个分区维护着一个事务信息的高水位日志类似于Kafka的offset机制。整体流程可以拆成五个阶段FindCoordinatorProducer根据自身TransactionalId哈希找到对应的Transaction CoordinatorInitTransactionsProducer向Coordinator请求写入事务标记获取一个PID并和这个Coordinator绑定关系BeginTransactionProducer在本地标记事务开始注意这个阶段不会通知CoordinatorProduce/Consume/FenceProducer发送消息到目标分区同时向Coordinator注册该分区参与事务Commit/AbortProducer主动提交或中止事务Coordinator将最终状态写入事务日志并在事务日志中标记所有参与分区为已提交或已中止。3.2 事务型Producer的完整代码姿态先看一段实际可运行的配置和代码Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, localhost:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); // 关键配置开启事务必须指定 transactional.id props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, tx-account-01); // 事务的超时时间 props.put(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG, 60000); // 开启幂等性Kafka 3.0默认开启 props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // 事务要求所有副本确认且不允许乱序 props.put(ProducerConfig.ACKS_CONFIG, all); KafkaProducerString, String producer new KafkaProducer(props); producer.initTransactions(); try { producer.beginTransaction(); producer.send(new ProducerRecord(topic-order, order-001, create)); producer.send(new ProducerRecord(topic-points, order-001, add-100)); producer.commitTransaction(); } catch (Exception e) { producer.abortTransaction(); }注意几个细节initTransactions()必须在线程启动时执行一次它会向Coordinator注册TransactionalId并申请PIDbeginTransaction()和commitTransaction()必须成对出现abortTransaction()用于异常回滚同一个Producer实例可以顺序开启多个事务但每个时刻只能有一个事务处于进行中如果并发调用多个线程共享同一个Producer并开启事务会抛出ConcurrentTransactions异常。3.3 consume-transform-produce模式与事务的配合事务最常见的落地场景是读一个Topic的数据处理之后写入另一个Topic。如果这个过程中间出现了异常最怕的是消息读出来了、处理完了但还没来得及写目标Topic下游消费者已经把不完整的数据消费走了。事务配合Consumer手动提交可以实现经典的EOS流转props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, read_committed); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); try (KafkaConsumerString, String consumer new KafkaConsumer(props)) { producer.initTransactions(); while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(200)); producer.beginTransaction(); for (ConsumerRecordString, String record : records) { producer.send(new ProducerRecord(topic-processed, record.key(), process(record.value()))); } // 把offset的提交也纳入事务 producer.sendOffsetsToTransaction( consumer.position(topic-input, 0) 0 ? currentOffsets : currentOffsets, consumer.groupMetadata()); producer.commitTransaction(); } }这里有个容易忽视的点producer.sendOffsetsToTransaction()会把Consumer的offset提交和消息写入放进同一个事务里。这样一来要么写入和提交同时成功要么同时失败回滚。即使处理流程中途崩溃重启后事务回滚offset也回到原位消息不会被重复处理。这就是Kafka事务在高版本里能做EOS的底气。4. 精确一次处理EOS的工程落地与边界很多人一听到精确一次就觉得这是Kafka的开箱功能但实际上EOS是一个端到端的语义需要生产端、Broker端、消费端三方配合才能成立。不同场景下EOS的落地方式差异很大。4.1 Kafka Streams中的EOS实现如果直接用Kafka Streams做流式处理EOS是相对完整地被框架封装好的。Kafka Streams内部通过以下机制协同Transactional Producer所有的中间结果和最终输出都通过事务写入State Store与Changelog Topic处理过程中的状态快照通过事务写入对应的changelog Topic消费位置与输出的一致性通过事务保证处理完某条记录、更新了状态、输出到下游这三个动作要么全部完成要么全部不完成。实际使用Kafka Streams时只需要在配置中设置props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);4.2 普通业务系统的伪精确一次方案但绝大多数业务系统并不是Kafka Streams而是消费消息 写数据库 调接口的传统架构。在这种场景下依赖Kafka事务并不能完全做到EOS因为事务只覆盖了Kafka内部的写入数据库写入和外部接口调用是事务无法感知的。我在实际项目中用过一套有效的组合方案思路是消费端把消息的唯一业务ID作为数据库主键或唯一索引重复插入直接报冲突业务逻辑忽略并提交offset外部接口调用前先查本地状态表如果状态已经是已处理直接跳过Kafka端开启read_committed隔离级别避免读到未提交的事务数据。这套方案的原理是依靠业务去重而不是依靠Kafka去重。它虽然不能叫严格意义上的EOS但对大多数订单、支付、业务流水场景来说效果和EOS一样而且实现成本低得多。4.3 EOS的真实边界和性能损耗EOS不是免费的午餐它的性能代价实打实存在。开启事务后Producer每次事务提交都要和Transaction Coordinator进行两轮RPCPrepareCommit和Commit并且每条消息都要附带事务标记Broker端还要多写一条控制消息。在低吞吐场景下这个代价不明显但在每秒数万条消息的高压力场景下吞吐量可能下降20%~30%。还有几个比较容易踩的配置坑配置项建议值原因transaction.state.log.replication.factor3事务日志必须保证高可用transaction.state.log.min.isr2至少2个ISR副本否则事务提交会失败transaction.max.timeout.ms默认为15分钟如果事务执行超过这个时间Coordinator会直接中止事务isolation.levelread_committed或read_uncommittedEOS场景必须设为read_committed5. 生产环境中的常见陷阱与排查经验最后这部分是我最想写的——理论谁都懂真正让一个系统崩掉的往往不是知识盲区而是一些看起来无关紧要的边界条件和配置音效。5.1 幂等性冲突与序列号溢出问题幂等性有一个很少有人提到的隐藏坑序列号溢出。PID的序列号是一个Int类型最大值为Int.MAX_VALUE。在一个长时间运行的Producer实例中如果单个分区写入速度特别快序列号总有一天会达到上限。遇到这种情况Broker端记录的ProducerState会选择重置或者触发异常。我处理的方案是在监控里加入对序列号增长速率的告警一旦接近阈值主动重启Producer实例以更换PID。另外一个常见的坑是Producer端的重试次数和幂等性冲突。虽然开启了幂等性后重复写入会被Broker识别并丢弃但如果retries设置得太大并且分区本身不可用Producer会一直卡在重试队列里导致消息延迟被无限放大。5.2 事务超时与Coordinator切换后的恢复流程我在生产环境遇到过一次非常诡异的问题某个消息始终无法成功消费查看日志发现Consumer端不断报TransactionCoordinator not available。排查链路是这样的先看Coordinator是否正常kafka-transactions.sh describe发现事务状态是PrepareCommit但一直卡着不动再看事务超时配置发现业务方设置了transaction.timeout.ms36000001小时远超Broker端的transaction.max.timeout.ms默认值15分钟Coordinator判定超时后自动中止了事务最终结果是Producer端毫不知情还在继续尝试提交Coordinator却已经开始回滚那些写入的消息。这个案例给了两个教训事务超时时间必须在Producer端和Broker端同时配置并且Producer端的值不能大于Broker端的最大值另外当Coordinator发生切换比如Broker重启时Producer需要重新执行initTransactions()来和新的Coordinator建立联系否则事务状态会一直停留在旧Coordinator上。5.3 一个真实的重复消费排查案例去年有一个订单相关的服务上线后对账程序发现订单量和支付流水量对不上多出了0.3%的重复支付回调。排查过程打开Consumer日志确认 offset提交正常情况下理论上不会重复消费进一步观察发现服务每次发布重启后会有部分消息被重新消费。原因是消费者组发生了Rebalance而Rebalance发生在offset提交之后的窗口期新分区的消费者从旧offset重新拉取数据启用Kafka事务之后配合sendOffsetsToTransaction把offset提交和业务事务绑定重复问题才被根治。这个案例给我最大的启发是重复消费本质上是处理完成和offset提交这两个动作不在同一个原子范围内。只靠Kafka的幂等性或者只靠业务去重表都能解决一部分问题但真正稳妥的方案是把消费-处理-提交作为一个整体来设计。最后说说我的体会算起来从第一次在生产环境配置acksall到现在我在Kafka消息保证这条路上踩过的坑、翻过的车可以写满一页纸。现在我对消息系统的态度就一句话不要迷信任何单一机制把幂等性、事务、业务去重表三者组合好才能真正做到精确一次。Kafka的幂等性解决的是Broker层的重复写入事务解决的是跨分区和跨Topic的原子性而端到端的精确一次处理永远是业务系统设计消息系统能力数据存储约束三方协作的结果。如果你正在做一个对消息一致性要求很高的系统我建议先从消费端的业务去重入手再逐步引入Kafka事务而不是一开始就把所有希望押在某一层的魔法开关上。
返回列表