惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

雷峰网
雷峰网
月光博客
月光博客
S
Security Affairs
宝玉的分享
宝玉的分享
D
DataBreaches.Net
MongoDB | Blog
MongoDB | Blog
Cloudbric
Cloudbric
PCI Perspectives
PCI Perspectives
B
Blog RSS Feed
腾讯CDC
Application and Cybersecurity Blog
Application and Cybersecurity Blog
Attack and Defense Labs
Attack and Defense Labs
N
News and Events Feed by Topic
T
The Blog of Author Tim Ferriss
H
Help Net Security
Vercel News
Vercel News
W
WeLiveSecurity
U
Unit 42
S
SegmentFault 最新的问题
Microsoft Azure Blog
Microsoft Azure Blog
Google Online Security Blog
Google Online Security Blog
云风的 BLOG
云风的 BLOG
Google DeepMind News
Google DeepMind News
S
Schneier on Security
The Register - Security
The Register - Security
酷 壳 – CoolShell
酷 壳 – CoolShell
Recent Announcements
Recent Announcements
博客园 - Franky
H
Hacker News: Front Page
WordPress大学
WordPress大学
I
Intezer
M
MIT News - Artificial intelligence
博客园 - 叶小钗
The Last Watchdog
The Last Watchdog
T
Troy Hunt's Blog
Stack Overflow Blog
Stack Overflow Blog
Microsoft Security Blog
Microsoft Security Blog
L
Lohrmann on Cybersecurity
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
I
InfoQ
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
Schneier on Security
Schneier on Security
P
Proofpoint News Feed
V
V2EX
Help Net Security
Help Net Security
小众软件
小众软件
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
大猫的无限游戏
大猫的无限游戏
The GitHub Blog
The GitHub Blog
F
Fortinet All Blogs

博客园_首页

Plist 二进制格式 Milvus 和 PGVector,哪个更好? OpenClaw 已过时?在 VS Code 中运行 Hermes Agent! 第30篇文章:一个大三计科生的自白 Manim如何在数学公式中完美显示中文? Docker 部署 RocketMQ 5 并发编程核心概念辨析 C#事务处理最佳实践:别再让“主表存了、明细丢了”的破事发生 CLI 是什么?为什么大厂突然集体卷命令行? 【从0到1构建一个ClaudeAgent】协作-自主Agent UIImageView 设置图片不生效的原因排查 最小二乘问题详解20:无先验约束下的增量式SFM自由网平差 痞子衡嵌入式:大话双核i.MXRT1180之XIP应用里借助MU实现可靠Flash IAP的方法 AI Chat 封装, SemanticKerne.AiProvider.Unified 已发布 Windows下右键编辑js文件无法打开记事本——在注册表中使用环境变量 在后台服务中使用 Scoped 服务,为什么总是报错? H200 安装驱动并使用sglang启动模型 wireshark 抓包Trap上报告警内容 我用 AI 辅助开发了一系列小工具(2):图片压缩工具 [A Primer On MC and CC] 2.1 Memory Consistency 1 - 指令重排序和 SC 模型 Oracle数据库SCN推进技术详解与实践指南 玩转控件:封装个带图片的Label控件 Claude Code 4.7 真正该升级的不是模型,而是你的工作流 前端小白一句话,AI 帮我做了个颜值拉满的桌面媒体播放器。当代码不再是门槛,一句话编程就是现实。 5. WorkBuddy: 小龙虾的灵魂三件套,让你的小龙虾不只是工具 SQLite 分片方案实战:三种分片策略的深度对比 告别简陋 UI!一款基于 Fluent Design 和基于 WinUI 的开源免费、现代化的 Avalonia UI 控件库 关于二进制排列组合枚举的总结 AI开发-python-LangGraph框架(3-27-LangGraph从零实现大模型智能决策工作流) ElasticSearch主分片和副本分片概念详解 【002】HTTPS 粗解:证书、TLS 握手与对后端配置的影响 Hermes Agent 一周暴涨五万 Star,但我劝你别急着追 明明连接的是Redis的DB0,为什么能查到DB3的数据? 【从0到1构建一个ClaudeAgent】协作-Agent团队 熟悉电子元器件之后,电子小白下一步该怎么走? MAF快速入门(23)通过C#类定义Skills .NET 高级开发 | 手写一个对象映射框架 FastAPI数据库ORM怎么选?我肝了三个Demo后,终于不再纠结了 mysqldump 参数拾遗:在遗忘与铭记之间 C# .NET 周刊|2026年3月5期 Claude code入门 - 陈彦斌 一文学习入门 ThingsBoard 开源物联网平台 GitHub 热门项目 | 2026年04月16日 如何为GIT设置全局勾子,为每次提交追加信息 Number.isFinite和isFinite与isNaN()和Number.isNaN的区别 PortSwigger SQL注入LAB2 推荐一个测试人必备的Skills,从功能到性能全搞定(附详细实操和安装下载方式) 筑基期:掌握Odoo基础核心知识点02(Odoo XML 开发方式详解) GLM模型这么火,咱们用vllm也咧一个呗! 深入理解 AbortController:从底层原理到跨语言设计哲学 字符串学习笔记 多租户系统框架的基础模块设计和分析设计 Apache SeaTunnel Zeta 为什么能做到“又快又稳”? AI开发-python-LangGraph框架(3-26-LangGraph基本概念及第一个简单样例) Vue 3 组件通信,别只会用 Props 和 Emits 了,这几个狠活儿你得看看 ElasticSearch7.X版本配置密码 用Manim实现动态交点计算--从一个动点问题说起 团结引擎+Addressable+Instant Game打包抖音小游戏 function call 实战:让 LLM 自动判断 pod 异常、调用日志工具并完成故障分析 bubseek —— 让 Agent 的足迹,变成团队的洞察 通过 C# 读取并导出 PDF 书签 如何用 GitHub Actions 实现 Steam 自动化发布 【从0到1构建一个ClaudeAgent】并发-后台任务 .NET 高级开发 | 定制 ASP.NET Core 框架 电子小白:什么是运算放大器(运放) zero2Agent:面向大厂面试的 Agent 工程教程,从概念到生产的完整学习路线 堆上的ORW HC32F460 USB CDC通信异常:非对齐访问异常排查 20260413-Hyperbridge 攻击事件:发生在默克尔山上的验证绕过 那些喊着AI 要淘汰你的人,正在靠你的焦虑赚大钱! 深度学习进阶(八)Swin Transformer 最小二乘问题详解19:带先验约束的增量式SFM优化与实现 SnapTranslate 3.0 正式发布:全局划词翻译 + 完整英语学习闭环,一站式搞定查词、记词、复习 工作的意义、工作的困难认知再思考 .NET + AI 进阶实战:基于类的技能开发 - 打造可治理的 Agent 能力模块 【从0到1构建一个ClaudeAgent】规划与协调-技能 上周热点回顾(4.6-4.12) 电子小白的工具三件套:面包板、杜邦线、万能板 单表五亿数据的查询优化 | Mysql、StarRocks 2. WorkBuddy:从“我是谁”到“帮我干活” C# 如何减少代码运行时间:7 个实战技巧 基于HelixToolkit.SharpDX 渲染3D模型 - 笺上知微 从零开始的双臂具身VLA起源及现阶段发展综述 - SkyXZ 记对 xonsh shell 的使用, 脚本编写, 迁移及调优 - pluvium27 受够了Vibe Coding的失控?换个起点,让AI事半功倍 从开始配置漏洞环境到漏洞复现流程 - 難しい 关于10年工作经验的程序员对OpenClaw的实战经验分享以及看法 - 虚无境 Any metadata 的内存布局 C# .NET 周刊|2026年3月2期 - InCerry 我帮你测过了,测试圈排名第二的 Skill 依然很牛逼 Skill Discovery | 无监督技能发现的经典工作总结 - MoonOut 上下文工程是什么?过时了么?一文讲明白! - 一枫说码 开了 TUN 模式还是直连?90% 的人都踩过这个坑 AScript扩展多种脚本语言 - rockey627 AI 学习笔记:Agent 的记忆机制 你能被装进一个文件里吗?——7 万人把同事"蒸馏"成了 AI - 我没有三颗心脏 Claude Code 通关手册(七):给 AI 装上技能包——Skills 完全指南 - 暮色之狐 在浏览器中快速编辑代码:VSCode Web 集成实践 - Newbe36524 蒸馏自己 skill?基于 Deepseek 的蒸馏器,丐版蒸馏方式,简单便捷 - To_Carpe_Diem Spring AI Aliababa和AgentScope,哪个更好? - 苏三说技术
从零学习Kafka:幂等与事务
Jackeyzhe · 2026-05-12 · via 博客园_首页

前文我们聊了生产者压缩机制相关的知识点,压缩机制主要是为了节省资源并提高吞吐。本文我们再来看一下 Kafka 的可靠性机制。

写在前面

在正式开始之前,先来回顾一下前面留的一个问题:Consumer 在解压数据时,需要关系 Producer 使用了哪种压缩算法吗?

答案是不需要,在 Kafka 消息批次的头信息中,包含有 Attribute 字段,这个字段中有几位专门用来标识这个批次所使用的压缩算法。Consumer 在拿到数据后,如果发现是压缩数据,会先使用对应的压缩算法解压,然后把解压后的数据交给业务代码。

可靠性保证

在实时计算领域,可靠性保证通常都有三种:

  • At most once:最多一次,即消息可能丢失,但不会重复。

  • At least once:最少一次,即消息不会丢失,但可能重复。

  • Exactly once:精确一次,消息不会丢失,也不会重复。

Kafka 默认支持的可靠性保证是第二种——最少一次。其原理也很简单,就是 Producer 端在没有收到 Broker 的确认消息时,会进行重试。支持最多一次的方法也很简单,只需要禁用重试就可以了,这样 Producer 写入消息可能失败也可能成功,但不会重复。

在对于消息准确性要求高的场景中,最多一次和最少一次都不能满足要求,需要保证精确一次才可以。在 Flink 中是通过 Checkpoint 来保证精确一次的,而在 Kafka 中是通过幂等和事务来保证的。

幂等性

什么是幂等性

我们先来看幂等性,首先解释一下什么是幂等性。在数学中的定义是,对于一个函数 f(x) 作用在任意一个元素上一次和多次,得到的结果是一样的。例如乘以 1 这个操作,一个数乘以一次 1 和乘以 n 次 1,得到的结果都相同,乘以 1 这个操作就是幂等的。而加 1 就不是。

在计算机领域,幂等性的定义是用户对同一个操作发起一次或多次请求,最终资源的状态是不变的。这就保证了我们可以放心的重试,不需要担心数据重复。

Kafka的幂等性

我们知道了幂等性的定义之后,重点来关注一下 Kafka 是如何开启幂等性的。默认情况下,Producer 不是幂等性的,我们可以通过设置一个参数来创建幂等性的 Producer。

props.put("enable.idempotence", true);
// 或者
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);

那么 Kafka 是如何实现幂等性的呢?答案是通过引入两个关键标识:Producer ID 和 Sequence Number。

  • Producer ID:当一个启用了幂等性的 Producer 启动时,Broker 会为它分配一个唯一的 ID,这个 ID 在 Producer 生命周期内保持不变。

  • Sequence Number:Producer 为发往每个分区的每条消息都附带一个从 0 开始的递增的序列号。

有了这两个标识之后,Broker 会在内存中维护一个 (Producer ID, Sequence Number) 组合的高水位线,即 last_sequence_number。在消息写入过程中可能会遇到以下三种情况:

  1. 正常写入:收到的 Sequence Number 比 last_sequence_number 大 1,Broker 接收消息并更新 last_sequence_number。

  2. 重复消息:如果收到的 Sequence Number 小于或等于 last_sequence_number,那么代表这个消息之前已经收到过了,Broker 会丢弃消息,但会返回给 Producer 写入成功。

  3. 丢消息/乱序:如果 Sequence Number 比 last_sequence_number 大了不止 1,说明中间漏了消息,可能是迟到或者缺失,此时 Broker 会报错。

聊到这里,你是不是觉得启用了幂等性的 Producer 之后,貌似就能保证消息不重不丢了。前面我们提到了 Producer ID 在 Producer 生命周期内保持不变,如果 Producer 重启了,那么 Broker 会为它分配一个新的 Producer ID,这样就无法保证消息不重复了。幂等性的另一个局限是跨分区的原子性无法保证。针对这两个问题 Kafka 引入了事务。

事务

从 0.11 版本开始,Kafka 提供了对事务的支持。关于事务的定义我们就不过多介绍了,就是 ACID。在 Kafka 中就是解决跨分区的原子性和 Producer 重启可能导致的消息重复。

开启事务需要进行两项配置,同时要修改发送消息的代码。

  • 开启幂等性,设置 enable.idempotence=true

  • 设置 Producer 的 transaction.id,最好设置成一个有意义的名字

发送消息的代码需要增加事务相关的处理

producer.initTransactions();
try { 
    producer.beginTransaction();
    producer.send(record1);
    producer.send(record2);
    producer.commitTransaction();
} catch (KafkaException e) { 
    producer.abortTransaction();
}

事务流程

为了支持事务,Kafka 引入了一个核心组件事务协调器和一个内部 Topic __transaction_state。整个过程基本就是两阶段提交的过程。

1. 初始化

Producer 启动时,会向 Broker 发送请求,确定属于自己的事务协调器,然后告诉事务协调器自己的 transaction.id。事务协调器收到请求之后,会记录 transaction.id 的 epoch,这是一个递增的批次号,有了这个批次号之后,旧的生产者就会被屏蔽。

2. 开启事务

初始化完成之后,需要开启事务,即调用 beginTransaction() 方法。这是在告诉事务协调器,我要开始往 Topic1 的分区 0 发送消息了。事务协调器就会把这些信息记录在__transaction_state 这个 Topic 中。

3. 发送消息

此时 Producer 开始发送消息,这些消息会写入 Broker 的磁盘,但它们会带有一个“事务进行中”的标识。

4. 提交阶段

发送完消息后,生产者调用 commitTransaction() 方法来提交事务,这里有两个步骤。

  1. 准备阶段:事务协调器将事务状态改为 PrepareCommit 并写入 __transaction_state

  2. 事务协调器向所有的分区 Leader 发送 WriteTxnMarker 请求,各个分区 Leader 会写入一个控制消息(Control Marker),如果中途有失败,事务协调器会发送 Abort 标记。

5. 完成

当所有的分区的 Marker 都写入完成时,事务协调器将事务状态改为 Complete,否则标记为 Abort。

至此,一次完整的事务消息发送的过程就结束了。

如何Exactly Once

前面的过程中,不管事务成功与否,消息都被写入到了 Broker 的磁盘中了,那如何保证 Exactly Once 呢?实际上在 Consumer 端还有一个参数控制:isolation.level,它用来控制隔离级别。这个参数取值有两个:

  • read_uncommitted:这是默认的,读未提交。也就是说,不管消息是什么状态,只要写到了 Broker 的磁盘上,都能被 Consumer 读到。

  • read_committed:这是读已提交,也是我们开启事务时需要的。它表明我们只能看到提交的事务消息。

这样幂等性+事务+隔离级别三者配合起来,Kafka 就能支持精确一次的可靠性保证了。

总结

我们来总结一下今天的内容,首先介绍了 Kafka 保证的可靠性,默认是最少一次,可以支持最多一次和精确一次。精确一次的支持依赖于幂等性、事务和隔离级别。我们从这三个方面分别介绍了在 Kafka 中是如何支持的。

最后我们留下一个问题。如果生产者已经发送了大量消息,但在最后提交之前突然宕机,事务协调器会如何处理这个未完成的事务呢?