news 2026/8/13 10:43:59

头疼的 Kafka 消息重复问题,从根上解决!

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
头疼的 Kafka 消息重复问题,从根上解决!

一、前言

数据重复这个问题其实也是挺正常,全链路都有可能会导致数据重复。

通常,消息消费时候都会设置一定重试次数来避免网络波动造成的影响,同时带来副作用是可能出现消息重复。

整理下消息重复的几个场景:
  1. 生产端:遇到异常,基本解决措施都是重试

    • 场景一:leader分区不可用了,抛LeaderNotAvailableException异常,等待选出新leader分区。

    • 场景二:Controller所在Broker挂了,抛NotControllerException异常,等待Controller重新选举。

    • 场景三:网络异常、断网、网络分区、丢包等,抛NetworkException异常,等待网络恢复。

  2. 消费端:poll一批数据,处理完毕还没提交offset,机子宕机重启了,又会poll上批数据,再度消费就造成了消息重复。

怎么解决?

先来了解下消息的三种投递语义:

  • 最多一次(at most once):消息只发一次,消息可能会丢失,但绝不会被重复发送。例如:mqttQoS = 0

  • 至少一次(at least once):消息至少发一次,消息不会丢失,但有可能被重复发送。例如:mqttQoS = 1

  • 精确一次(exactly once):消息精确发一次,消息不会丢失,也不会被重复发送。例如:mqttQoS = 2

了解了这三种语义,再来看如何解决消息重复,即如何实现精准一次,可分为三种方法:
  1. Kafka幂等性Producer保证生产端发送消息幂等。局限性,是只能保证单分区且单会话(重启后就算新会话)

  2. Kafka事务:保证生产端发送消息幂等。解决幂等Producer的局限性。

  3. 消费端幂等: 保证消费端接收消息幂等。蔸底方案。

1)Kafka幂等性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. 注意
  1. 设置幂等,启动幂等。

  2. 配置acks,注意:一定要设置acks=all,否则会抛异常。

  3. 配置max.in.flight.requests.per.connection需要<= 5,否则会抛异常OutOfOrderSequenceException

    • 0.11 >= Kafka < 1.1,max.in.flight.request.per.connection = 1

    • Kafka >= 1.1,max.in.flight.request.per.connection <= 5

为了更好理解,需要了解下Kafka 幂等机制:

  1. Producer每次启动后,会向Broker申请一个全局唯一的pid。(重启后pid会变化,这也是弊端之一)

  2. Sequence Numbe:针对每个<Topic, Partition>都对应一个从0开始单调递增的Sequence,同时Broker端会缓存这个seq num

  3. 判断是否重复:<pid, seq num>Broker里对应的队列ProducerStateEntry.Queue(默认队列长度为 5)查询是否存在

    • 如果nextSeq == lastSeq + 1,即服务端seq + 1 == 生产传入seq,则接收。

    • 如果nextSeq == 0 && lastSeq == Int.MaxValue,即刚初始化,也接收。

    • 反之,要么重复,要么丢消息,均拒绝。

这种设计针对解决了两个问题:
  1. 消息重复:场景Broker保存消息后还没发送ack就宕机了,这时候Producer就会重试,这就造成消息重复。

  2. 消息乱序:避免场景,前一条消息发送失败而其后一条发送成功,前一条消息重试后成功,造成的消息乱序。

那什么时候该使用幂等:
  1. 如果已经使用acks=all,使用幂等也可以。

  2. 如果已经使用acks=0或者acks=1,说明你的系统追求高性能,对数据一致性要求不高。不要使用幂等。

2)Kafka事务

使用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 Producer<String, String> producer = new KafkaProducer<String, String>(props); // 初始化事务 producer.initTransactions(); try{ // 开始事务 producer.beginTransaction(); // 发送数据 producer.send(new ProducerRecord<String, 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 a86dff3689b68f6af7eb3da5a21c2dba06e9623f3c961154a8bbbe3e9991dea4

2、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 Producer<String, String> producer; final String outTopic; public KafkaProducerApplication(final Producer<String, 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 ProducerRecord<String, 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 Producer<String, 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 { List<String> 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.

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/13 10:42:29

深度解析报纸门户网站建设方案:从传统媒体转型到数字化生存的实战指南,助力媒体融合新跨越

在这个信息爆炸、碎片化阅读盛行的时代,如果你还认为建一个报纸门户网站只是把报纸上的文字搬到网上,那可就大错特错了。现在的媒体环境,早已不是“我写你看”的单向传播时代,而是双向互动、全媒体矩阵、数据驱动的新生态。很多老牌报社或者新兴的资讯平台在搞数字化转型时…

作者头像 李华
网站建设 2026/8/13 10:42:39

4.1.2三目运算符

作用&#xff1a;通过三目运算符莱判断 语法&#xff1a;表达式1 ? 表达式2 : 表达式3 示例&#xff1a; #include<iostream> using namespace std;int main() {//三目运算符//创建三个变量 a b c//将a和b作比较&#xff0c;将变量大的值赋值给变量cint a 10;int b 20…

作者头像 李华
网站建设 2026/8/13 10:39:10

网站建设合同编号全攻略:如何通过正规流程规避建站陷阱并保障企业权益

在这个数字化浪潮席卷全球的时代,每一个企业,无论你是开在闹市区的一家传统杂货铺,还是活跃在硅谷的创新科技公司,拥有一个像样的网站已经不再是“锦上添花”,而是“生存必需”。然而,当老板拍着胸脯说“我们要搞个大动作,做个高端网站”的时候,作为执行层面的项目负责…

作者头像 李华
网站建设 2026/8/13 10:36:10

machine 框格标注的特殊规定(形位公差)

T、公共公差带&#xff08;核心概念&#xff09;公共公差带多个零件表面 / 轴线&#xff0c;功能上需要共用同一个公差带统一约束&#xff0c;统一管控平面度、直线度等几何精度&#xff0c;就叫公共公差带。最典型两类应用&#xff1a;共面&#xff08;多个端面齐平&#xff0…

作者头像 李华
网站建设 2026/8/13 10:35:10

红外多类别无人机 YOLOv11 检测 基于 YOLOv11n 的红外多类型无人机目标检测系统 智慧红外飞行检测 - 红外多类别无人机航拍数据集

基于 YOLOv11n 的红外多类型无人机目标检测系统 智慧红外飞行检测 智慧红外飞行检测-红外多类别无人机目标检测数据集&#xff0c;6271张&#xff0c;yolo&#xff0c;voc&#xff0c;coco三种标注方式 图像尺寸:640*640 类别数量:2类 训练集图像数量:4390; 验证集图像数量:1…

作者头像 李华
网站建设 2026/8/13 10:34:42

揭秘杭州91网站建设背后的真实故事:如何打造一台真正懂生意的营销利器

咱们今天不整那些虚头巴脑的互联网黑话,也不扯什么元宇宙、区块链的大道理,就踏踏实实聊点接地气的东西。在咱们杭州,提起“杭州91网站建设”,很多老板第一反应可能是一个网址,或者是一行代码,但在行内人眼里,它其实是一场关于“信任”和“转化”的持久战。我在这个行业…

作者头像 李华