解决NATS JetStream多消费者实例监听Work Queue报错问题
NATS JetStream Work Queue多消费者实例启动报错解决
问题场景
我已创建一个NATS JetStream发布者,通过以下方式向“Work Queue”流发布消息:
import nats import json import asyncio from pathlib import Path NATS_URL = "nats://localhost:4222" dummy_tasks = [ ('traffic.download.URMO.v1.1', 1), ('traffic.download.URMO.v1.2', 2), ('traffic.download.URMO.v1.3', 3), ('traffic.download.NUMO.v1.1', 4), ('traffic.download.NUMO.v1.2', 5), ('traffic.download.URMO.v1.4', 6), ('traffic.download.URMO.v1.5', 7), ('traffic.download.NUMO.v1.3', 8), ] async def main(): # Create a work queue nc = await nats.connect(NATS_URL) js = nc.jetstream() stream_config = nats.js.api.StreamConfig( name="traffic-download", subjects=["traffic.download.*.*.*"], retention=nats.js.api.RetentionPolicy.WORK_QUEUE, ) await js.add_stream(stream_config) for task in dummy_tasks: msg_payload = json.dumps(task[1]).encode('utf-8') await js.publish(task[0], msg_payload) print(f"Sent message for task {task[0]}") await asyncio.sleep(1) await nc.close() if __name__ == '__main__': asyncio.run(main())
目标是让多个进程实例监听该流并“获取”任务进行处理,但目前仅能运行一个消费者实例,运行多个时会遇到错误。消费者代码如下:
async def callback(msg): #payload = json.loads(msg.data.decode('utf-8')) #write_geojson(payload['path'], payload['bounds']) print(f"Received message for task {msg.subject}") async def main(): nc = await nats.connect(NATS_URL) js = nc.jetstream() await js.subscribe("traffic.download.*.*.*", stream="traffic-download", cb=callback) await asyncio.Event().wait() if __name__ == '__main__': asyncio.run(main())
报错信息:
nats.js.errors.BadRequestError: nats: BadRequestError: code=400 err_code=10100 description='filtered consumer not unique on workqueue stream'
问题原因
Work Queue类型的流不允许存在多个重复的临时过滤消费者。当前的订阅方式没有指定队列名称,会创建临时的过滤消费者,当第二个消费者实例启动时,NATS会检测到重复的过滤规则,因此抛出错误。
解决方案
要实现多消费者协作处理Work Queue任务,需要使用队列订阅(Queue Subscription),即给消费者指定一个统一的队列名称,让多个消费者加入同一个队列组。
修改后的消费者代码:
async def callback(msg): #payload = json.loads(msg.data.decode('utf-8')) #write_geojson(payload['path'], payload['bounds']) print(f"Received message for task {msg.subject}") async def main(): nc = await nats.connect(NATS_URL) js = nc.jetstream() # 指定队列名称,多个消费者使用同一个队列名称即可加入队列组 await js.subscribe( "traffic.download.*.*.*", stream="traffic-download", cb=callback, queue="task-processing-queue" # 新增队列名称参数 ) await asyncio.Event().wait() if __name__ == '__main__': asyncio.run(main())
说明
- 指定
queue参数后,NATS会将所有使用该队列名称的消费者归为一个队列组,Work Queue流中的消息会被均匀分发给队列组内的各个消费者,每个消息只会被一个消费者处理。 - 这种方式完全符合Work Queue的设计逻辑:任务被消费确认后从流中移除,多个消费者并行处理任务,实现负载均衡。
内容的提问来源于stack exchange,提问作者Jamess11
相关产品推荐
相关产品推荐

