ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Kafka 幂等生产者与事务消息在大促交易防重中的性能损耗评估

Kafka 幂等生产者与事务消息在大促交易防重中的性能损耗评估 Kafka 幂等生产者与事务消息在大促交易防重中的性能损耗评估在大促交易与资金结算链路中“消息绝不能丢但也绝不能重复消费”是一条不可动摇的资金安全红线。如果用户支付成功的一条 MQ 消息被重复消费了两次导致用户钱包被双重扣款或商家被重复打款会直接引发严重的资损和客诉。Kafka 默认提供的传输语义是“至少一次At-Least-Once”——当生产端发送消息成功但由于网络抖动未收到 Broker 的 ACK 时生产端会发起重试从而在 Topic 中产生两条内容完全一致的重复消息或者当消费端处理完毕但未能成功提交 Offset 时Rebalance 后也会导致消息被再次投递。为了实现精确一次Exactly-Once Semantics, EOSKafka 官方推出了幂等生产者Idempotent Producer与事务消息Transactional Messaging。然而天下没有免费的午餐。在大促数十万 TPS 极限压测下开启这些特性究竟会带来多少性能损耗架构师该如何进行技术选型与权衡幂等生产者Idempotent Producer的实现机理与损耗从 Kafka 3.0 开始生产者默认开启了幂等性enable.idempotencetrue。底层工作原理Producer IDPID分配生产者在初始化时会向 Broker 申请一个全局唯一的 64 位整数 PID序列号自增Sequence Number生产者发送给每个Topic, Partition的每一批消息都会附带一个单调递增的 Sequence Number从 0 开始Broker 端内存滑动窗口去重Broker 在内存中为每个PID, Partition维护最近 5 个已提交消息的 Sequence Number 状态。当 Broker 收到一条 Sequence Number $\le$ 当前已落盘最大序号的消息时Broker 直接丢弃该重复消息并正常向生产者返回成功 ACK从而在单个分区内部彻底杜绝了网络重试引发的数据重复。[Producer] --(PID: 1001, Seq: 42, Data: OrderPaid)-- [Broker Partition-0] | 检查内存窗口: 当前最大 Seq41 - 成功写入更新 Seq42 | [Producer] --(重试发送: PID: 1001, Seq: 42)--------- [Broker Partition-0] | 检查内存窗口: 已存在 Seq42 - 丢弃 Payload返回 SUCCESS ACK!压测性能损耗评估在 16 核 32G 机器、单实例并发发送 50,000 TPS 的基准压测下延迟LatencyP99 延迟从 4.2ms 增加至 4.4ms增幅仅为4.7%吞吐Throughput由于每条消息仅增加了 PID8 字节和 Sequence Number4 字节共 12 字节的微小开销Broker 端仅需进行一次内存哈希比对整体吞吐量下降小于3%核心结论幂等生产者的开销微乎其微大促全链路所有核心与非核心 Topic 均建议强制开启Kafka 事务消息Transactional API的实现机理与代价幂等生产者只能保证单个生产者在单个分区内部的防重且无法保证“跨多个 Topic 生产与消费的原子性”。在大促订单与库存协同场景下通常需要满足原子操作“消费订单创建消息 $\rightarrow$ 本地数据库更新 $\rightarrow$ 向库存扣减 Topic 发送新消息”三者必须要么全部成功要么全部回滚。为此Kafka 引入了基于两阶段提交2PC的事务协调器Transaction Coordinator机制。// 典型的 Kafka 事务生产者代码支持跨分区原子写入 KafkaProducerString, String producer new KafkaProducer(props); producer.initTransactions(); // 向 Transaction Coordinator 注册全局 transactional.id try { producer.beginTransaction(); producer.send(new ProducerRecord(trade-order-topic, orderJson)); producer.send(new ProducerRecord(stock-deduct-topic, stockJson)); // 提交事务协调器向事务日志写入 COMMIT 标记并向各分区写入 Control Batch (Commit Marker) producer.commitTransaction(); } catch (ProducerFencedException | OutOfOrderSequenceException e) { producer.close(); } catch (KafkaException e) { producer.abortTransaction(); // 异常回滚 }事务消息的底层物理开销多轮跨网络 RPC 往返开启事务后单次提交必须经历AddPartitionsToTxn、EndTxn等多次与 Transaction Coordinator专门的__transaction_state内部 Topic的同步网络交互控制标记块Control Batch写入每个参与事务的分区都必须额外写入一条特殊的 Commit/Abort 控制标记段导致底层磁盘 I/O 写入次数翻倍消费端延迟阻塞Read Committed消费端必须配置isolation.levelread_committed消费者拉取数据时如果某个事务尚未提交后续的所有正常消息全部无法被消费直到该事务结束LSO 推进极易引发长尾延迟抖动。压测性能损耗数据在相同硬件规格与并发载荷下延迟Latency单次事务提交的 P99 延迟从 4.2ms 恶化至18.6ms膨胀了 4.4 倍吞吐Throughput单 Broker 最大并发写入 TPS 从 85,000 跌落至62,000吞吐骤降约 27%。大促最佳实践业务去重表 轻量幂等生产的黄金组合综合评估大促峰值的高吞吐诉求与资损防范我们强烈推荐如下分层组合架构[生产端] - 开启 Kafka 原生轻量幂等性 (enable.idempotencetrue) - 几乎零性能损耗 | v (网络传输) [消费端] - 采用【本地 Redis 预占 数据库防重唯一索引表】实现业务级 Exactly-Once放弃重型的 Kafka 分布式事务除非是在流计算如 Flink 端到端 EOS场景在常规的 Java OLTP 微服务中坚决避免使用重量级的 Kafka Transactional API防止其 27% 的吞吐损耗拖垮网关消费端落地“防重唯一索引”消费端在写入数据库业务表的同时将message_id全局唯一分布式雪花 ID写入一张包含UNIQUE KEY (msg_id)的trade_dedup_log去重表利用 MySQL 本地事务的原子性一旦发生重复消费数据库底层唯一键冲突报错抛出DuplicateKeyException消费者直接安全 ACK既保障了 100% 不重复又榨干了 Kafka 的极致吞吐。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进