使用正则订阅Pulsar新主题时消息投递延迟的原因排查
Pulsar正则订阅新主题消息延迟45秒问题排查
问题描述
我有一个Python程序,用于接收Pulsar消息并更新文档数据库。当消息发布至已存在的主题时,运行正常;但当发布时创建新主题,订阅接收该新主题的消息会延迟约45秒。
消费者代码
def event_received(consumer, message): event = Event.from_json(message.properties()['event_type'], message.data().decode('utf-8')) history = load_topic(f'persistent://demo/person/{event.id}') person = Person('{"id": "%s"}' % event.id) person.load(history) person.apply(event) collection.replace_one({"id": event.id}, person.to_map(), True) aggregates.subscribe(re.compile('demo/person/.*'), 'person-mongodb', message_listener=event_received, regex_subscription_mode=pulsar.RegexSubscriptionMode.PersistentOnly, initial_position=pulsar.InitialPosition.Earliest) while True: try: pass except KeyboardInterrupt: aggregates.close()
生产者代码
app.post('/person/created') def person_created(): producer = aggregates.create_producer(f'persistent://demo/person/{request.json["id"]}') producer.send(json.dumps(request.json).encode('utf-8'), properties={'event_type': 'PersonCreated'}) producer.close() return {}
原因分析
核心原因是Pulsar Broker的正则订阅主题发现机制存在默认扫描间隔。当使用正则表达式订阅主题时,Broker不会实时检测新创建的匹配主题,而是按照默认的45秒时间间隔去扫描匹配的主题列表。新创建的主题必须等到下一次扫描周期完成后,才会被加入到订阅的主题集合中,此时消费者才能开始接收该主题的消息,这就导致了45秒左右的延迟。
解决方案
- 调整Broker主题扫描间隔:修改Broker配置文件(如
broker.conf)中的brokerTopicsPatternCheckIntervalSeconds参数,将其设置为更小的值(例如5秒),让Broker更频繁地扫描匹配正则的新主题,从而缩短延迟。修改后需要重启Broker生效。 - 预先创建主题:如果业务场景允许,在生产者发送消息前,通过Pulsar Admin API或CLI提前创建好目标主题,这样正则订阅的消费者能立即识别到主题,不会产生延迟。
- 优化主题设计:避免为每个用户ID创建独立主题,改用单主题+分区、或通过消息属性过滤的方式,从根源上减少新主题的创建,同时规避正则订阅的扫描延迟问题。
- 检查自动主题创建配置:确保Broker开启了自动创建主题功能(默认开启),相关配置如
allowAutoTopicCreation保持启用,但此配置对正则订阅的扫描间隔无直接影响,主要还是依赖扫描间隔参数的调整。
内容的提问来源于stack exchange,提问作者Brian Richardson
相关产品推荐
相关产品推荐

