在Pulsar中如何排除消息发送方对应的Topic消费者
多会话场景下避免消息重复推送的解决方案
核心问题分析
当前架构中,向user:A主题发送消息时,所有订阅该主题的客户端会话(包括发起消息的会话#1)都会收到消息,导致会话#1重复获取已经通过REST响应拿到的内容。这个问题无需重构整体架构,通过消息标识过滤或主题路由优化即可解决。
可行解决方案
1. 消息携带排除会话ID,会话Consumer过滤
- 处理用户发送消息的REST请求时,从请求上下文(如Cookie、请求头中的会话标识)获取发起请求的会话ID(如
session#1)。 - 向Pulsar的
user:A主题发送消息时,在消息payload或元数据中添加exclude_session_id字段,值为发起会话的ID。 - 每个会话的Consumer收到消息后,先对比自身会话ID与
exclude_session_id:若一致则直接丢弃消息,不推送给对应WebSocket客户端;若不一致则正常推送。 - 优势:改动极小,仅需在消息生产和消费环节增加简单逻辑,无需调整Pulsar配置或架构核心。
2. 拆分用户主题为会话专属主题,通过Hub路由消息
- 为每个WebSocket会话创建专属Pulsar主题,命名规则如
user:{userId}-{sessionId}(例如user:A-session#1、user:A-session#2)。 - 服务端完成消息存储后,调用Hub接口获取目标用户(如用户A)的所有活跃会话列表,排除发起消息的会话#1。
- 向剩余会话对应的专属主题分别发送消息,每个会话的Consumer仅订阅自己的专属主题,自然不会收到针对其他会话的消息。
- 优势:消息推送精准,无需客户端过滤逻辑;缺点是需要Hub维护会话与主题的映射关系,服务端发送消息时需多一步路由判断。
3. 利用Pulsar订阅属性实现消息过滤
- 每个会话的Consumer订阅
user:A主题时,设置订阅属性session_id: "{sessionId}"(如session_id: "session#1")。 - 生产消息时,使用Pulsar的消息过滤功能,指定消息不匹配带有
session_id = "session#1"属性的订阅。 - 优势:过滤逻辑在Pulsar层面完成,客户端和服务端无需额外处理;需确认Pulsar版本是否支持基于订阅属性的消息过滤(Pulsar 2.8+支持相关功能)。
内容的提问来源于stack exchange,提问作者BVtp
相关产品推荐
相关产品推荐

