头疼的 Kafka 消息重复问题,从根上解决!
一、前言数据重复这个问题其实也是挺正常全链路都有可能会导致数据重复。通常消息消费时候都会设置一定重试次数来避免网络波动造成的影响同时带来副作用是可能出现消息重复。整理下消息重复的几个场景生产端遇到异常基本解决措施都是重试。场景一leader分区不可用了抛LeaderNotAvailableException异常等待选出新leader分区。场景二Controller所在Broker挂了抛NotControllerException异常等待Controller重新选举。场景三网络异常、断网、网络分区、丢包等抛NetworkException异常等待网络恢复。消费端poll一批数据处理完毕还没提交offset机子宕机重启了又会poll上批数据再度消费就造成了消息重复。怎么解决先来了解下消息的三种投递语义最多一次at most once消息只发一次消息可能会丢失但绝不会被重复发送。例如mqtt中QoS 0。至少一次at least once消息至少发一次消息不会丢失但有可能被重复发送。例如mqtt中QoS 1精确一次exactly once消息精确发一次消息不会丢失也不会被重复发送。例如mqtt中QoS 2。了解了这三种语义再来看如何解决消息重复即如何实现精准一次可分为三种方法Kafka幂等性Producer保证生产端发送消息幂等。局限性是只能保证单分区且单会话重启后就算新会话Kafka事务保证生产端发送消息幂等。解决幂等Producer的局限性。消费端幂等 保证消费端接收消息幂等。蔸底方案。1Kafka幂等性Producer幂等性指无论执行多少次同样的运算结果都是相同的。即一条命令任意多次执行所产生的影响均与一次执行的影响相同。幂等性使用示例在生产端添加对应配置即可Properties props new Properties(); props.put(enable.idempotence, ture); // 1. 设置幂等 props.put(acks, all); // 2. 当 enable.idempotence 为 true这里默认为 all props.put(max.in.flight.requests.per.connection, 5); // 3. 注意设置幂等启动幂等。配置acks注意一定要设置acksall否则会抛异常。配置max.in.flight.requests.per.connection需要 5否则会抛异常OutOfOrderSequenceException。0.11 Kafka 1.1,max.in.flight.request.per.connection 1Kafka 1.1,max.in.flight.request.per.connection 5为了更好理解需要了解下Kafka 幂等机制Producer每次启动后会向Broker申请一个全局唯一的pid。重启后pid会变化这也是弊端之一Sequence Numbe针对每个Topic, Partition都对应一个从0开始单调递增的Sequence同时Broker端会缓存这个seq num判断是否重复拿pid, seq num去Broker里对应的队列ProducerStateEntry.Queue默认队列长度为 5查询是否存在如果nextSeq lastSeq 1即服务端seq 1 生产传入seq则接收。如果nextSeq 0 lastSeq Int.MaxValue即刚初始化也接收。反之要么重复要么丢消息均拒绝。这种设计针对解决了两个问题消息重复场景Broker保存消息后还没发送ack就宕机了这时候Producer就会重试这就造成消息重复。消息乱序避免场景前一条消息发送失败而其后一条发送成功前一条消息重试后成功造成的消息乱序。那什么时候该使用幂等如果已经使用acksall使用幂等也可以。如果已经使用acks0或者acks1说明你的系统追求高性能对数据一致性要求不高。不要使用幂等。2Kafka事务使用Kafka事务解决幂等的弊端单会话且单分区幂等。Tips这块篇幅较长这先稍微提及下使用之后另起一篇。事务使用示例分为生产端 和 消费端Properties props new Properties(); props.put(enable.idempotence, ture); // 1. 设置幂等 props.put(acks, all); // 2. 当 enable.idempotence 为 true这里默认为 all props.put(max.in.flight.requests.per.connection, 5); // 3. 最大等待数 props.put(transactional.id, my-transactional-id); // 4. 设定事务 id ProducerString, String producer new KafkaProducerString, String(props); // 初始化事务 producer.initTransactions(); try{ // 开始事务 producer.beginTransaction(); // 发送数据 producer.send(new ProducerRecordString, String(Topic, Key, Value)); // 数据发送及 Offset 发送均成功的情况下提交事务 producer.commitTransaction(); } catch (ProducerFencedException | OutOfOrderSequenceException | AuthorizationException e) { // 数据发送或者 Offset 发送出现异常时终止事务 producer.abortTransaction(); } finally { // 关闭 Producer 和 Consumer producer.close(); consumer.close(); }这里消费端Consumer需要设置下配置isolation.level参数read_uncommitted这是默认值表明Consumer能够读取到Kafka写入的任何消息不论事务型Producer提交事务还是终止事务其写入的消息都可以读取。如果你用了事务型Producer那么对应的Consumer就不要使用这个值。read_committed表明Consumer只会读取事务型Producer成功提交事务写入的消息。当然了它也能看到非事务型Producer写入的所有消息。3消费端幂等“如何解决消息重复” 这个问题其实换一种说法就是如何解决消费端幂等性问题。只要消费端具备了幂等性那么重复消费消息的问题也就解决了。典型的方案是使用消息表来去重上述栗子中消费端拉取到一条消息后开启事务将消息Id新增到本地消息表中同时更新订单信息。如果消息重复则新增操作insert会异常同时触发事务回滚。二、案例Kafka 幂等性 Producer 使用环境搭建可参考https://developer.confluent.io/tutorials/message-ordering/kafka.html#view-all-records-in-the-topic准备工作如下1、Zookeeper本地使用Docker启动$ docker run -d --name zookeeper -p 2181:2181 zookeeper a86dff3689b68f6af7eb3da5a21c2dba06e9623f3c961154a8bbbe3e9991dea42、Kafka版本2.7.1源码编译启动看上文源码搭建启动3、启动生产者Kafka源码中exmaple中4、启动消息者可以用Kafka提供的脚本# 举个栗子topic 需要自己去修改 $ cd ./kafka-2.7.1-src/bin $ ./kafka-console-producer.sh --broker-list localhost:9092 --topic test_topic创建topic1副本2 分区$ ./kafka-topics.sh --bootstrap-server localhost:9092 --topic myTopic --create --replication-factor 1 --partitions 2 # 查看 $ ./kafka-topics.sh --bootstrap-server broker:9092 --topic myTopic --describe生产者代码public class KafkaProducerApplication { private final ProducerString, String producer; final String outTopic; public KafkaProducerApplication(final ProducerString, String producer, final String topic) { this.producer producer; outTopic topic; } public void produce(final String message) { final String[] parts message.split(-); final String key, value; if (parts.length 1) { key parts[0]; value parts[1]; } else { key null; value parts[0]; } final ProducerRecordString, String producerRecord new ProducerRecord(outTopic, key, value); producer.send(producerRecord, (recordMetadata, e) - { if(e ! null) { e.printStackTrace(); } else { System.out.println(key/value key / value \twritten to topic[partition] recordMetadata.topic() [ recordMetadata.partition() ] at offset recordMetadata.offset()); } } ); } public void shutdown() { producer.close(); } public static void main(String[] args) { final Properties props new Properties(); props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); props.put(ProducerConfig.ACKS_CONFIG, all); props.put(ProducerConfig.CLIENT_ID_CONFIG, myApp); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); final String topic myTopic; final ProducerString, String producer new KafkaProducer(props); final KafkaProducerApplication producerApp new KafkaProducerApplication(producer, topic); String filePath /home/donald/Documents/Code/Source/kafka-2.7.1-src/examples/src/main/java/kafka/examples/input.txt; try { ListString linesToProduce Files.readAllLines(Paths.get(filePath)); linesToProduce.stream().filter(l - !l.trim().isEmpty()) .forEach(producerApp::produce); System.out.println(Offsets and timestamps committed in batch from filePath); } catch (IOException e) { System.err.printf(Error reading file %s due to %s %n, filePath, e); } finally { producerApp.shutdown(); } } }启动生产者后控制台输出如下启动消费者$ ./kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic myTopic修改配置 acks启用幂等的情况下调整acks配置生产者启动后结果是怎样的修改配置acks 1修改配置acks 0会直接报错Exception in thread main org.apache.kafka.common.config.ConfigException: Must set acks to all in order to use the idempotent producer. Otherwise we cannot guarantee idempotence.修改配置 max.in.flight.requests.per.connection启用幂等的情况下调整此配置结果是怎样的将max.in.flight.requests.per.connection 5会怎样当然会报错Caused by: org.apache.kafka.common.config.ConfigException: Must set max.in.flight.requests.per.connection to at most 5 to use the idempotent producer.

相关新闻

本地AI部署实战:从环境搭建到API集成的完整指南

本地AI部署实战:从环境搭建到API集成的完整指南

这次我们来看一个名为“重返未来great idea”的项目。从名称上看,它可能是一个结合了创意、未来感或某种特定主题的生成式AI工具或应用。这类项目通常聚焦于通过AI技术实现特定风格的图像、视频或文本生成,其核心价值在于能否在本地环境稳定运行&#xf…

2026/9/21 18:05:34 阅读更多 →
元初混沌体系架构 第二卷 第三十六篇 7G时空通信稳态闭环总复盘

元初混沌体系架构 第二卷 第三十六篇 7G时空通信稳态闭环总复盘

第三十六篇 7G时空通信稳态闭环总复盘单元收官总序 第二单元完整定位复盘本单元为第二卷 鸿蒙7G星际超域通信升维体系核心中层支柱单元,单元编号19–36,共计18篇正统学术篇章,精准承接第一单元本源公理体系,为后续频谱升维、抗干…

2026/9/22 12:53:08 阅读更多 →
BeyondCompare:专业源代码比对工具的核心功能与实战应用

BeyondCompare:专业源代码比对工具的核心功能与实战应用

1. 项目概述:为什么我们需要专业的源代码比对工具?在软件开发、代码审计、版本合并或者团队协作的日常里,一个场景反复出现:你手上有两份看起来相似但又不同的代码文件,可能是不同分支的修改,可能是重构前后…

2026/9/22 4:31:33 阅读更多 →

最新新闻

图解Enclave原理:微服务升级踩坑实录

图解Enclave原理:微服务升级踩坑实录

图解Enclave原理:微服务升级踩坑实录 昨天凌晨三点,生产环境报警炸了。 版本升级后 API 全变了,之前跑得好好的 Enclave 服务,这次直接报错。 我盯着屏幕上的 ECS Exception…

2026/9/22 12:52:40 阅读更多 →
3个坑搞定AccessPoint调试,Go语言最佳实践

3个坑搞定AccessPoint调试,Go语言最佳实践

3个坑搞定AccessPoint调试,Go语言最佳实践 复制来的 AccessPoint 代码跑不通,报错信息模糊,改一行崩一行?别慌。这是很多后端开发者接手旧项目或参考 GitHub…

2026/9/22 12:52:40 阅读更多 →
软文是啥?转岗开发必看的速查手册

软文是啥?转岗开发必看的速查手册

软文是啥?转岗开发必看的速查手册 刚转岗做开发,是不是觉得手里全是零散的语法知识,却拼不出一个完整的项目?很多人卡在“懂代码”到“能落地”这一步,急需一份 速查手册 来理清思路。今天不聊虚的,直接拆解一个让无数新人头秃的隐性成本——…

2026/9/22 12:52:40 阅读更多 →
2026最新Nyan Cat项目配置避坑:5个报错一次讲透

2026最新Nyan Cat项目配置避坑:5个报错一次讲透

2026最新Nyan Cat项目配置避坑:5个报错一次讲透 刚接手那个老项目的同事,是不是也被 Nyan Cat 这个前端特效卡得怀疑人生?明明只是加个彩虹猫跑马灯,结果 npm install 还没跑完, webpack 直接报…

2026/9/22 12:51:39 阅读更多 →
遥感信息处理避坑指南:3个完整示例搞定API变更

遥感信息处理避坑指南:3个完整示例搞定API变更

遥感信息处理避坑指南:3个完整示例搞定API变更 版本升级后 API 全变了,是不是让你抓狂?刚写好的脚本跑不起来,报错信息看得头大。别慌,我整理了遥感信息处理的完整示例,帮你快速上手。…

2026/9/22 12:51:39 阅读更多 →
5步搞定无限的未知win7性能瓶颈,实战项目提速3倍

5步搞定无限的未知win7性能瓶颈,实战项目提速3倍

5步搞定无限的未知win7性能瓶颈,实战项目提速3倍 官方文档翻了三遍还是晕?别慌,很多老手都卡在这。无限的未知win7这种底层机制,光看理论根本跑不起来。拿一个 实战项目 实测,你才会发现哪里在拖后腿。…

2026/9/22 12:51:39 阅读更多 →

日新闻

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天

3台商务办公笔记本实测:手写实现环境配置,告别卡半天 配置环境就卡半天?别怪机器慢,多半是你没选对工具链。在Java、Go或Python的项目现场, 手写实现…

2026/9/22 0:00:41 阅读更多 →
剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑

剑帝加点速查手册:3分钟搞懂核心逻辑 面试被问原理答不上来,是不是常态?别慌。很多开发者对着 GitHub 开源仓库里的代码发呆,看似简单实则暗藏玄机。今天这份【剑帝加点】速查手册,直接带你拆解核心实现,把面试必考的原理讲透。…

2026/9/22 0:00:41 阅读更多 →
手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优

手写实现图片压缩网站核心:搞定WebP转换与质量调优 复制来的代码跑不通不知道怎么调?别慌,这种“复制粘贴地狱”在开发圈太常见了。尤其是做 图片压缩网站…

2026/9/22 0:00:41 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/22 4:32:41 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/22 4:38:57 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/22 8:51:04 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/21 15:36:51 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/21 15:36:51 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/22 2:43:42 阅读更多 →