Kafka如何避免消息重复消费
Producer端
重复生产场景
Producer的send()方法可能会出现异常,配合生产者参数retries>0,生产者会在出现可恢复异常的时候进行重试。
若出现不可恢复异常的时候,配合send()的异步发送方式,则可能在回调函数中进行消息重发。上述均可能导致消息重复。
解决方法
Kafka的幂等性就是为了避免出现生产者重试的时候出现重复写入消息的情况。
开启幂等性功能配置(该配置默认为false)如下:
prop.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG,true);
Consumer端
重复消费场景
一、自动提交消费位移
kafka默认消费位移的提交是自动提交,由消费者参数enable.auto.commit配置,默认为true。
这个自动提交并不是每消费一条消息就自动提交消费位移,而是定期提交,这个定期提交的时间由客户端参数auto.commit.interval.ms配置,默认5秒。
下一次就还得在上一次消费位移的位置重新开始消费,造成重复消费:
-
如果在拉取消息进行消费,但是下一次提交位移之前消费者崩溃了。 在消费者关闭之前调用了consumer.unsubscribe()方法取消订阅。
解决办法
设置手动提交消费位移。
二、Kafka服务端的Partition再均衡机制导致消息重复消费
在Kafka中有一个Partition Balance机制,就是把多个Partition均衡的分配给多个消费者。消费端会从分配到的Partition里面去消费消息,如果消费者在默认的5分钟内没有处理完这一批消息。就会触发Kafka的Rebalance机制,从而导致offset自动提交失败。而Rebalance之后,消费者还是会从之前没提交的offset位置开始消费,从而导致消息重复消费。
解决办法
提高消费端的处理性能避免触发Balance,比如可以用多线程的方式来处理消息,缩短单个消息消费的时长。或者还可以调整消息处理的超时时间,也还可以减少一次性从Broker上拉取数据的条数。
