Apache Pulsar消息去重工作机制及配置后仍收重复消息问题咨询
结论
你遇到的情况属于预期行为。Apache Pulsar的消息去重机制不是基于消息内容的去重,和AWS SQS的content based deduplication逻辑完全不同。
Pulsar去重机制的核心逻辑
Pulsar的去重是为了解决生产者重试场景下的重复消息问题,属于基于序列号的幂等性去重,判断消息重复的三个核心维度是:
- 固定的生产者名称(
producerName) - 消息对应的序列号(
sequence ID) - 消息所属的主题
只有三个维度完全匹配的消息,Broker才会判定为重复并自动丢弃。
你当前代码的问题
你的代码中虽然10次发送的消息内容完全相同,但你没有手动指定消息的sequence_id,Pulsar客户端会在每次调用send()方法时自动生成递增的唯一序列号,因此10条消息的序列号互不相同,Broker不会判定为重复,消费者自然会收到全部10条消息。
如果要验证Pulsar原生去重能力,你可以修改send逻辑手动指定相同的序列号:
import pulsar import json client = pulsar.Client('pulsar://localhost:6650') producer = client.create_producer( 'persistent://public/default/my-topic', send_timeout_millis=0, producer_name="producer-1") data = {'key1': 0, 'key2' : 1} for i in range(10): encoded_data = json.dumps(data).encode('utf-8') # 手动指定相同的序列号,触发Pulsar去重 producer.send(encoded_data, sequence_id=1) client.close()
修改后运行代码,消费者只会收到1条消息。
若需要实现基于内容的去重的可选方案
如果你的业务需要类似SQS的按内容去重能力,可以自行实现逻辑:
- 生产侧对消息内容做哈希计算,将哈希值作为消息的
sequence_id或者消息key,相同内容的消息会拿到相同的哈希值,即可触发Pulsar原生去重 - 消费侧消费时基于内容哈希做幂等校验
- 利用Pulsar Functions做流处理层的内容去重
内容的提问来源于stack exchange,提问作者Shubham Jain
相关产品推荐
相关产品推荐

