如何在Pulsar中实现同Key消息的自定义顺序管控与过期操作过滤?
针对Pulsar同Key消息版本冲突的解决方案
1. Key-Shared订阅+消息属性过滤(推荐)
通过Pulsar原生的Key-Shared订阅模式,结合消息中携带的版本ID,在消费阶段自动丢弃旧版本的无效消息:
- 配置Consumer为
Key_Shared订阅类型,确保同一Key的消息只会被单个Consumer实例处理,避免并发处理导致的版本判断混乱。 - 发送消息时,在消息属性中携带两个关键信息:
versionId(用来证明操作先后的版本号,#1的版本号大于#2)、operation(标记是CREATE还是DELETE)。 - 初始化Consumer时设置
MessageFilter,维护每个Key的已处理最大版本ID:- 收到DELETE消息时,更新该Key的最大版本ID并正常处理。
- 收到CREATE消息时,若其版本ID小于该Key已记录的最大版本ID,直接丢弃;否则正常处理。
示例Java代码:
// 线程安全的Map,记录每个Key的已处理最大版本ID ConcurrentHashMap<String, Long> keyMaxVersionMap = new ConcurrentHashMap<>(); Consumer<String> consumer = pulsarClient.newConsumer(Schema.STRING) .topic("your-target-topic") .subscriptionName("version-aware-sub") .subscriptionType(SubscriptionType.Key_Shared) .messageFilter((consumer, msg) -> { String key = msg.getKey(); String op = msg.getProperty("operation"); long currentVersion = Long.parseLong(msg.getProperty("versionId")); // 处理DELETE操作,更新最大版本 if ("DELETE".equalsIgnoreCase(op)) { keyMaxVersionMap.put(key, currentVersion); return true; } // 处理CREATE操作,仅当版本大于已记录的最大版本才放行 Long maxVersion = keyMaxVersionMap.get(key); return maxVersion == null || currentVersion > maxVersion; }) .subscribe();
2. 结构化Schema+版本校验
如果使用Avro、JSON等结构化Schema,可直接在Schema中定义versionId和operation字段,逻辑和上述方案一致:
- 在Schema里添加
versionId(长整型)和operation(枚举类型)字段。 - 同样采用Key-Shared订阅,Consumer处理消息时提取Schema中的版本信息,对比后过滤旧版本CREATE消息。
3. 自定义Key级SequenceId
Pulsar内置SequenceId是全局的,但你可以在消息属性中自定义keySequenceId(和Key绑定的序列值),替代全局SequenceId实现Key维度的顺序控制,配合上述过滤逻辑,同样能达到丢弃旧消息的效果。
所有方案均基于Pulsar原生能力实现,无需额外自定义复杂流程,核心是通过Key-Shared订阅保证同Key消息的单实例处理,再利用版本ID判断消息有效性。
内容的提问来源于stack exchange,提问作者Brian Lee
相关产品推荐
相关产品推荐

