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

推荐订阅源

T
Tailwind CSS Blog
V
Vulnerabilities – Threatpost
Hacker News - Newest:
Hacker News - Newest: "LLM"
V2EX - 技术
V2EX - 技术
N
News and Events Feed by Topic
cs.AI updates on arXiv.org
cs.AI updates on arXiv.org
F
Fortinet All Blogs
C
Cybersecurity and Infrastructure Security Agency CISA
A
Arctic Wolf
美团技术团队
Cloudbric
Cloudbric
大猫的无限游戏
大猫的无限游戏
WordPress大学
WordPress大学
S
Security Affairs
T
Tenable Blog
MyScale Blog
MyScale Blog
W
WeLiveSecurity
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
aimingoo的专栏
aimingoo的专栏
博客园_首页
K
KPMG report finds enterprise disconnect between AI and its ROI | CIO
月光博客
月光博客
宝玉的分享
宝玉的分享
博客园 - Franky
PCI Perspectives
PCI Perspectives
H
Hacker News: Front Page
D
Docker
博客园 - 司徒正美
L
LINUX DO - 最新话题
GbyAI
GbyAI
A
About on SuperTechFans
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
I
InfoQ
Microsoft Azure Blog
Microsoft Azure Blog
Google DeepMind News
Google DeepMind News
Microsoft Security Blog
Microsoft Security Blog
T
Threatpost
Google DeepMind News
Google DeepMind News
Help Net Security
Help Net Security
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
人人都是产品经理
人人都是产品经理
C
Check Point Blog
MongoDB | Blog
MongoDB | Blog
Y
Y Combinator Blog
AI
AI
有赞技术团队
有赞技术团队
N
News | PayPal Newsroom
T
The Blog of Author Tim Ferriss
D
DataBreaches.Net
O
OpenAI News

博客园_首页

Linux实操--组管理、权限管理和定时任务 Java + EasyExcel 实现单个接口导出多个Excel Mem0 源码解析系列(二):提示词工程的深度剖析 Openclaw TaskFlow究竟是什么?和普通Skill技能有什么区别 博文阅读密码验证 - 博客园 嘉立创开源:应该是全网MicroPython教程最多的开发板 Hermes Agent 集成实践:从协议到生产 2026年AI编程工具横评:Cursor、Codex、Claude Code、Zed、Windsurf Java程序员必看的RAG入门教程 2026 AI效率神器:Superpowers + Claude Code 保姆级教程 本地大模型部署全攻略:从 0 到 1 玩转 Ollama 【从0到1构建一个ClaudeAgent】内存管理-上下文压缩 .NET 高级开发 | 设计、实现一个事件总线框架 电子小白入门之NE555 3. WorkBuddy:隐藏玩法,一键召唤专家,让 AI 以"专家身份"给你干活 和AI一起搞事情#3:Claude Teammate 游戏开发翻车实录 【OpenClaw】通过 Nanobot 源码学习架构---(7)Memory C# .NET 周刊|2026年3月3期 我在 Debian 11 上把 K8s 单机搭起来了,过程没你想的那么顺(/opt 目录版) 深度学习进阶(七)Data-efficient Image Transformer CLI+Skill搭建浏览器AI自动化框架,告别一切重复枯燥任务 告别Token账单无底洞:OpenClaw本地部署,重塑企业数据主权的唯一解 FastAPI+Vue:文件分片上传+秒传+断点续传,这坑我帮你踩平了! SBTI 爆火后,我做了个程序员版的 CBTI。。已开源 + 附开发过程 多模态检索开始进入工程期:用 Sentence Transformers 搭建可落地的 Multimodal RAG 100多行代码实现一个最简单的Agent(用ReAct) Claude Code 通关手册(八):推荐 5 个 Hooks,代码质量提升 3 倍 老板:“有人截图了!”。安全部门:“收到,马上查暗水印!” - why技术 技术之外,皆是人间 C#/.NET/.NET Core技术前沿周刊 | 第 69 期(2026年4.01-4.12) Snack JSONPath 项目架构分析 Claude Code Buddy 小析:一个非核心功能,如何体现产品的细节完成度 AI新时代下的图床管理方案-Cloudflare图床+MCP+Skills方案指南 化繁为简:顺丰速运App如何通过 HarmonyOS SDK实现专业级空间测量 从零实现富文本编辑器#13-React非编辑节点的内容渲染 AI开发-python-langchain框架(3-23-OpenAI Functions风格Tool Calling智能助手) .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 PbootCMS 网站内容数量多导致访问慢?这些实用优化方案帮你提速! - 家兴网络技术工作室 上下文工程是什么?过时了么?一文讲明白! - 一枫说码 网站漏洞怎么发现并修复?一篇实用指南(附完整流程) - 家兴网络技术工作室 开了 TUN 模式还是直连?90% 的人都踩过这个坑 Github日报|2026年04月12日 - AI一族 AScript扩展多种脚本语言 - rockey627 AI 学习笔记:Agent 的记忆机制 你能被装进一个文件里吗?——7 万人把同事"蒸馏"成了 AI - 我没有三颗心脏 Claude Code 通关手册(七):给 AI 装上技能包——Skills 完全指南 - 暮色之狐 在浏览器中快速编辑代码:VSCode Web 集成实践 - Newbe36524 蒸馏自己 skill?基于 Deepseek 的蒸馏器,丐版蒸馏方式,简单便捷 - To_Carpe_Diem Spring AI Aliababa和AgentScope,哪个更好? - 苏三说技术 Etsy 把 1000 个 MySQL 分片迁进 Vitess:425TB 数据背后的真正问题不是性能,而是运维规模 MicroPython LVGL基础知识和概念:底层渲染与性能优化 - FreakStudio 数据库草图算法 Python 潮流周刊#146:CPython 引入 Rust 的进展 - 豌豆花下猫 最小生成树 - mofei1116 红日靶场七:从外网入口、容器逃逸到 AD 接管的完整利用链复盘 - YouDiscovered1t 分享四款开源且实用的 Kafka 管理工具 - 追逐时光者 vLLM 权重加载机制全解析:从挑战到理想架构 LCT 学习笔记 - ACehomoxue Avalonia UI 12.0.0 正式发布:架构演进和性能飞跃 - 张善友 当 AI Agent 把调用链拉长,延迟开始成为一门生意 conhost.exe 无法显示 U+2717 - 145a 太秀了,我把自己蒸馏成了 Skill!已开源 - 程序员鱼皮 ASP.NET Core 内存缓存实战:一篇搞懂该怎么配、怎么避坑 基于 Ghostty 带有分割标签页和为 Claude 编程设计的通知终端 - BugShare AI 焊死入口:教育的“操作系统级”重塑 - 郝hai 初级Java开发工程师使用sql脚本编写代码的过程是简单而且不糊涂 - CoderOilStation Claude Code通关手册(六):MCP协议完全指南 - 暮色之狐 边框灯光环绕动画特效实现指南 - Newbe36524 开源:子木蒸馏版的 SEO 审计工具 seo-audit-skill v1.0 我所理解的Python元模型 【从0到1构建一个ClaudeAgent】规划与协调-TodoWrite - 程序员Seven Claude 和 Codex 在审计 Skill 上性能差异探究 - ACai_sec AScript如何实现中文脚本引擎 - rockey627 【渗透测试】HTB Season10 Garfield 全过程wp - dynasty_chenzi Android 开发者为什么必须掌握 AI 能力?端侧视角下的技术变革 树状数组正确性证明 - AC-wyr 你的 AI 焦虑,可能比 AI 本身更危险——ATM 机没有消灭银行柜员,但恐慌消灭了你的判断力 - 我没有三颗心脏 一个拉胯的分库分表方案有多绝望?整个部门都在救火! - 冰河团队 动态规划入门必学之走方格问题 - Ofnoname PostgREST 与 PostgreSQL 角色权限配置全解析(生产级实践) - SheepDog1998 使用 UEFI 图形输出协议 GOP 在屏幕上显示图像的方法 - 阿源- Claude Code通关手册(五):组建你的AI专家团队,子代理系统 - 暮色之狐 一个程序员到架构师的催婚路之感悟(整整10年后的催婚相亲感悟) - MisterLip 用 Agent Skill 自动生成工作周报 - 赵康
从零学习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 中是如何支持的。

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