You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.07 11:36:05