
一、前言数据重复这个问题其实也是挺正常全链路都有可能会导致数据重复。通常消息消费时候都会设置一定重试次数来避免网络波动造成的影响同时带来副作用是可能出现消息重复。整理下消息重复的几个场景生产端遇到异常基本解决措施都是重试。场景一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.