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

解决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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 18:57:35