当前位置: 首页 > news >正文

Kafka 事务消息实现详解

Kafka 事务消息实现详解

Kafka 从 0.11 版本开始引入事务,配合幂等生产者,实现了Exactly-Once语义。


一、为什么需要事务

1.1 典型的"消费-处理-生产"场景

Consumer ──读取──► Topic A ──处理──► 写入 Topic B

如果处理成功但写入失败,或者写入成功但 offset 未提交,就会产生:

  • 重复消费:写完 B 后崩溃,重启后再次消费同一条 A 的消息
  • 消息丢失:处理完 offset 已提交,但 B 写入失败

1.2 事务解决了什么

Kafka 事务提供的是原子多分区写入 + 消费位移提交,保证:

要么全部成功:消息写入 B + offset 提交 要么全部失败:消息不写入 + offset 不回滚(保持原样)

不是传统数据库的 ACID 事务——不涉及多行回滚、隔离级别等概念。


二、核心概念

2.1 四个关键角色

┌──────────────────────────────────────────────────────────┐ │ │ │ ┌──────────┐ ┌──────────────────┐ ┌─────────────┐ │ │ │ Producer │ │ Transaction │ │ __transaction│ │ │ │ (事务ID) │ │ Coordinator │ │ _state │ │ │ │ │ │ (TC, Broker上) │ │ (内部Topic) │ │ │ └──────────┘ └──────────────────┘ └─────────────┘ │ │ │ │ │ │ │ │ 注册/提交 │ 持久化状态 │ │ │ ├────────────────►├─────────────────────►│ │ │ │ │ │ │ │ ┌──────────┐ │ │ │ │ │ Consumer │ │ │ │ │ │ (隔离级别)│ │ │ │ │ └──────────┘ │ │ │ │ │ └──────────────────────────────────────────────────────────┘
角色说明
Transactional Producer配置transactional.id的生产者
Transaction Coordinator (TC)运行在 Broker 上的事务协调器,每个transactional.id会 hash 到固定的 TC
__transaction_state内部 Topic(50 分区),TC 将事务状态持久化到这里
Consumer (isolation.level)消费者通过隔离级别控制是否读取未提交的事务消息

2.2 关键 ID

ID含义生命周期
transactional.id用户配置的稳定事务标识跨 Producer 重启不变
PID(Producer ID)Broker 分配的 64 位数字Producer 重启后重新分配,但transactional.id不变时会复用旧的 epoch 过期
Producer Epoch单调递增,区分同 PID 的不同实例每次 Producer 重新 initTransactions 时 +1
Transactional Sequence Number每个分区内的消息序号幂等生产的基础
transactional.id: "order-service-tx-001" ← 用户配置,稳定 ↓ initTransactions 请求 PID: 123456789, Epoch: 0 ← Broker 分配 ↓ Producer 重启,再次 initTransactions PID: 123456789, Epoch: 1 ← PID 不变,Epoch 递增

三、事务完整流程(时间线)

3.1 创建事务生产者

Propertiesprops=newProperties();props.put("bootstrap.servers","localhost:9092");props.put("transactional.id","order-tx-001");// ← 关键:稳定的事务 IDprops.put("enable.idempotence",true);// ← 事务自动开启幂等// 幂等要求:acks=all, retries>0, max.in.flight.requests.per.connection≤5KafkaProducer<String,String>producer=newKafkaProducer<>(props);

3.2 完整流程图

Producer Transaction Coordinator __transaction_state │ │ │ │ ① initTransactions() │ │ ├───────────────────────────────►│ FindCoordinator │ │ │ 分配 PID + Epoch │ │◄───────────────────────────────┤ │ │ ② beginTransaction() │ │ │ (仅本地标记,不发网络请求) │ │ │ │ │ │ ③ send(topic-A, msg) │ │ ├───────────────────────────────►│ 消息写入分区但标记为未提交 │ │ │ │ │ ④ send(topic-B, msg) │ │ ├───────────────────────────────►│ 同上 │ │ │ │ │ ⑤ sendOffsetsToTransaction() │ │ ├───────────────────────────────►│ 将消费位移也纳入事务 │ │ │ │ │ ⑥ commitTransaction() │ │ ├───────────────────────────────►│ 写入 PrepareCommit │ │ ├─────────────────────────────►│ │ │ 写入 Committed │ │ ├─────────────────────────────►│ │ │ 向涉及的分区 Leader 发送 │ │ │ TransactionMarker(COMMIT) │ │ │ │ │◄───────────────────────────────┤ 返回成功 │

3.3 Java 完整示例

// ==================== 事务生产者 ====================Propertiesprops=newProperties();props.put("bootstrap.servers","localhost:9092");props.put("transactional.id","order-tx-001");// enable.idempotence 在配置 transactional.id 时自动为 trueKafkaProducer<String,String>producer=newKafkaProducer<>(props);producer.initTransactions();// ① 初始化:注册 PIDtry{producer.beginTransaction();// ② 开启事务// ③ 发送业务消息producer.send(newProducerRecord<>("order-topic","order-123","created"));// ④ 将消费位移也纳入事务(原子绑定)Map<TopicPartition,OffsetAndMetadata>offsets=newHashMap<>();offsets.put(newTopicPartition("source-topic",0),newOffsetAndMetadata(100));producer.sendOffsetsToTransaction(offsets,"consumer-group-1");producer.commitTransaction();// ⑤ 提交}catch(ProducerFencedExceptione){// PID 被 epoch 更新的实例抢占,当前实例已"僵尸"producer.close();}catch(KafkaExceptione){producer.abortTransaction();// ⑥ 异常时回滚}

3.4 消费者端

PropertiesconsumerProps=newProperties();consumerProps.put("bootstrap.servers","localhost:9092");consumerProps.put("group.id","consumer-group-1");consumerProps.put("isolation.level","read_committed");// ← 关键配置// read_committed: 只读已提交的事务消息(默认是 read_uncommitted)KafkaConsumer<String,String>consumer=newKafkaConsumer<>(consumerProps);

四、内部实现原理

4.1 事务消息在分区中的物理形态

Partition 0 中的消息序列(逻辑视图): ┌──────────────────────────────────────────────────────────────┐ │ offset │ key │ value │ 事务标记 │ ├────────┼───────────┼─────────────┼───────────────────────────┤ │ 0 │ order-123 │ "created" │ (无) │ │ 1 │ order-456 │ "created" │ (无) │ │ 2 │ order-789 │ "created" │ ┌─ PID=100, Epoch=0 │ │ 3 │ order-789 │ "paid" │ │ 事务未提交 │ ← read_committed 读不到 │ 4 │ │ COMMIT │ └─ TransactionMarker │ ← 控制消息,消费者不可见 │ 5 │ order-999 │ "created" │ ┌─ PID=100, Epoch=0 │ │ 6 │ │ ABORT │ └─ TransactionMarker │ ← 回滚标记 │ 7 │ order-111 │ "created" │ (无) │ └──────────────────────────────────────────────────────────────┘

关键点:

  • 事务消息照常写入分区,和普通消息混在一起
  • 区别在于消息头部带有事务元信息(PID、Epoch、Sequence Number)
  • 事务结束时,TC 会在每个涉及的分区末尾追加一条TransactionMarker(控制消息)
  • Consumerread_committed模式读到 ABORT 标记后,会跳过该事务的所有消息

4.2 __transaction_state 内部 Topic

┌──────────────────────────┐ │ __transaction_state │ │ 分区数: 50 (固定) │ │ 副本数: 3 (建议) │ │ 压缩策略: compaction │ ├──────────────────────────┤ │ Key: transactional.id │ │ Value: 事务状态消息 │ │ - Empty / Ongoing │ │ - PrepareCommit │ │ - PrepareAbort │ │ - CompleteCommit │ │ - CompleteAbort │ └──────────────────────────┘

TC 启动时通过读取__transaction_state恢复所有处于Ongoing状态的事务,然后继续推进。

4.3 事务提交流程深入

Producer Transaction 分区 Leader │ Coordinator │ │ commitTransaction() │ │ ├───────────────────────►│ │ │ │ 1. 写 PrepareCommit │ │ │ 到 __transaction │ │ │ _state │ │ │ │ │ │ 2. 对各分区 Leader │ │ │ 发送 WriteTxnMarker│ │ ├──────────────────────►│ │ │ │ 3. 追加 COMMIT │ │ │ Control Message │ │◄──────────────────────┤ │ │ │ │ │ 4. 写 Committed 到 │ │ │ __transaction_state │ │◄───────────────────────┤ │ │ 返回成功 │ │

两阶段提交?不完全是。Kafka 事务是1.5 阶段提交——TC 先持久化PrepareCommit,然后并行发 Marker 到各分区,全部成功后再持久化Committed。如果 TC 在中间宕机,重启后从__transaction_state恢复,重试发 Marker。

4.4 事务回滚

producer.abortTransaction();

流程与提交对称:

  1. TC 写PrepareAbort__transaction_state
  2. 向各分区发送 ABORT TransactionMarker
  3. read_committed消费者读到 ABORT 后,跳过该事务的消息(仿佛从未发送)

五、Consumer 隔离级别

5.1 read_uncommitted(默认)

读到所有消息,包括未提交和已回滚的事务消息 → 可能读到最终被回滚的"脏"数据

5.2 read_committed

只读到已提交事务的消息 → 实现 Exactly-Once 的前提 Consumer 会在内存中维护一个"已中止事务"的缓存: ┌──────────────────────────────────┐ │ 正在消费 offset=5 │ │ 读到 ABORT TransactionMarker │ │ → 把 (PID=100, Epoch=0) 记入缓存 │ │ → 回溯跳过该事务的所有消息 │ │ (offset 2, 3) │ │ → 继续从 offset=7 消费 │ └──────────────────────────────────┘

六、幂等生产者 (Idempotent Producer) —— 事务的前提

事务依赖于幂等生产,幂等又是独立可用的功能。

6.1 幂等原理

普通生产者: send(msg) → 网络超时 → 重试 → Broker 收到两条 msg(重复) 幂等生产者: send(msg, PID=100, Seq=5) → 网络超时 → 重试 send(msg, PID=100, Seq=5) Broker 收到第二条时: "PID=100 的 Seq=5 我已经有了,忽略" → 只保留一条

Broker 端每个分区维护:

PID → 最近 5 个 Seq 的去重窗口 ┌──────┬─────────────────────────┐ │ PID │ 已收到的 Seq 号 │ ├──────┼─────────────────────────┤ │ 100 │ [1, 2, 3, 4, 5] │ │ 101 │ [1, 2, 3] │ └──────┴─────────────────────────┘

6.2 幂等 vs 事务

对比幂等事务
开启条件enable.idempotence=truetransactional.id
作用域单分区内消息不重复跨分区原子写入
跨分区原子性
原子绑定消费位移✅(sendOffsetsToTransaction
PID 生成由 Broker 随机分配关联到transactional.id,可恢复

七、关键配置清单

Producer

# 事务 ID(设为非空即开启事务支持) transactional.id=my-app-tx-001 # 事务超时(默认 60 秒,最大 15 分钟) transaction.timeout.ms=60000 # 幂等(事务自动开启,无需显式设置) enable.idempotence=true # 幂等要求(自动设置) acks=all retries=2147483647 # Integer.MAX_VALUE max.in.flight.requests.per.connection=5

Consumer

# 隔离级别 isolation.level=read_committed # read_uncommitted | read_committed # 会话超时必须大于事务超时 session.timeout.ms > transaction.timeout.ms

Broker

# 事务状态日志副本数 transaction.state.log.replication.factor=3 # 事务状态日志分区数(默认 50,不可通过配置修改) transaction.state.log.num.partitions=50 # 事务状态日志分段大小 transaction.state.log.segment.bytes=104857600 # 事务超时最大值 transaction.max.timeout.ms=900000 # 15 分钟

八、僵尸实例隔离 —— Producer Fencing

场景: 1. Producer-A (Epoch=0) 开启事务,网络分区 2. Producer-A 被认为宕机 3. Producer-A 重启 → initTransactions → PID 复用,Epoch=1 4. 老的 Producer-A (Epoch=0) 网络恢复,尝试 commitTransaction Broker 收到 Epoch=0 的提交请求: "你的 Epoch 已经过期,当前 Epoch=1,拒绝!" → ProducerFencedException → 防止"僵尸"实例写入脏数据

这就是transactional.id必须稳定的原因:同一个 ID 恢复后 PID 不变但 Epoch 递增,旧的 Epoch 的所有操作自动失效。


九、适用场景与局限性

适用场景

场景示例
消费-转换-生产从 topic-A 读 → 处理 → 写 topic-B + 提交 offset
多分区原子写入订单创建同时写"订单主题"和"通知主题"
Kafka Streams内部大量使用事务保证 exactly-once

局限性

局限说明
不跨系统只能保证 Kafka 内部的原子性,不涉及数据库、Redis 等外部系统
性能开销事务提交有额外网络往返 +__transaction_state写入,吞吐下降约 10-20%
消费者需配合必须isolation.level=read_committed,且消费者会缓冲未提交消息
超时限制默认 60 秒,最长 15 分钟,不适合长事务
不覆盖 consumer 的非 Kafka 操作如果在消费-处理-生产中间写了数据库,数据库写入不受事务保护

十、常见问题

Q1: commitTransaction 没收到响应,是成功还是失败?

调用producer.commitTransaction()时如果超时或抛异常,不要重试 commit。正确做法是producer.close()然后重建 Producer。Broker 会自行完成或超时回滚。

Q2: 同一个 transactional.id 能多实例并发吗?

绝对不能。同一时刻只能有一个transactional.id的活跃 Producer。第二个initTransactions会触发 fencing,第一个被踢出。

Q3: 事务中的消息什么时候对消费者可见?

read_committed模式下,当 TransactionMarker(COMMIT)被追加到分区后,消费者才能读到该事务的消息。注意消费者可能滞后于 LSO(Last Stable Offset)。


十一、总结

Kafka 事务的本质: 幂等生产者 (PID + Seq) │ ├── 单分区不重复 │ ▼ 事务生产者 (PID + Epoch + transactional.id) │ ├── 跨分区原子写入 ├── 原子绑定消费位移 ├── 僵尸实例隔离 (Fencing) │ ▼ 实现 Exactly-Once 语义 (Kafka Streams EOS)

一句话:Kafka 事务 = 幂等发送 + 原子多分区提交 + 消费位移原子绑定,通过 Transaction Coordinator 和__transaction_state内部 Topic 协调,消费者配合read_committed实现端到端的 Exactly-Once。

http://www.cnnetsun.cn/news/3959733.html

相关文章:

  • 技术内容创作模式切换:从教程到研究写作的实践指南
  • SpringBoot+Vue构建心理健康测评系统:从架构设计到工程实践
  • 本地化媒体处理工具搭建:从视频分析到自动化剪辑的工程实践
  • Windows 10/11 通过 WSL 2 安装 Hadoop 3.1.3 单机环境完整指南
  • 抖音无水印下载神器:douyin-downloader 完全使用手册
  • Qt 实时曲线卡顿优化:从QPainter到OpenGL的3级加速实战
  • C++从重复代码到标准库:模板、STL与string入门
  • Simulink实现两区域电力系统二次调频与AGC控制
  • RAID 5配置全流程详解:从原理到实战的存储基石搭建
  • Unity集成海康威视RTSP视频流:基于UMP插件的跨平台监控方案
  • Elasticsearch核心架构与实战:从倒排索引到生产部署
  • 高效文件管理:从根目录批量处理到自动化工作流实践
  • Selenium无头浏览器实战:从原理到生产环境部署与优化
  • Win10系统光盘刻录全攻略:从镜像获取到高可靠性刻录与验证
  • 网络排障实战:从协议原理到经典案例的9个关键场景解析
  • 《基于机器学习的中风风险预测模型研究》3(设计源文件+万字报告+讲解)(支持资料、图片参考_相关定制)_文章底部可以扫码
  • LlamaIndex ResponseSynthesizer 详解:从检索到生成的 RAG 核心组件
  • LiDAR技术深度解析:从核心原理到工程实践全链路指南
  • 锐丰专业音频功率放大器G350风扇配件参数
  • MediaPipe+Unity实时动作捕捉:低成本实现3D角色驱动
  • CSP-J网络连接模拟题解析:字符串处理与状态管理实战技巧
  • 卷积神经网络(CNN)结构详解:从核心原理到工程实践
  • 动态稀疏注意力DSA:突破多模态大模型推理瓶颈的关键技术
  • 从零构建卷积神经网络:PyTorch实战CIFAR-10图像分类
  • Python实现凯撒密码:从古典密码到现代编程实践
  • 大模型选型实战指南:从榜单排名到场景落地的四维评估法
  • AI下半场_03_CSDN版_Token经济学
  • 2026年上海企业新闻发布资源平台哪家好?深度剖析及优选指南
  • CSS3实现缺角矩形、折角边框与折角效果:clip-path与渐变实战指南
  • Codex接入团队后,真正卡壳的不是写代码