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

NATS Server无法向多进程并行分发任务问题求助

NATS JetStream任务队列无法并行处理任务

背景

尝试用NATS JetStream替代Celery+RabbitMQ作为任务队列,编写了生产者和消费者脚本,但多消费者进程无法并行处理任务。

生产者脚本(producer.py)

import json  # 原代码遗漏该导入,需补充
import nats
from nats.js.api import RetentionPolicy

NATS_HOST = '0.0.0.0'
NATS_PORT = '4222'
DOCUMENT_EXT_SUBJECT = "file_process"

nc = await nats.connect(servers=f"nats://{NATS_HOST}:{NATS_PORT}")
js = nc.jetstream()
await js.add_stream(
    name="sample-stream", 
    subjects=[DOCUMENT_EXT_SUBJECT],
    retention=RetentionPolicy.WORK_QUEUE,
)

for i in range(20):
    ack = await js.publish(
        subject=DOCUMENT_EXT_SUBJECT, 
        payload=json.dumps(
            {
                "some_load" : "lsome_load", 
            }).encode()
    )

消费者脚本(consumer.py)

import asyncio
import signal
import sys
import logging

import nats
from nats.js.api import RetentionPolicy, ConsumerConfig, AckPolicy

consumer_config = ConsumerConfig(
    ack_wait=900,
    max_deliver=1, 
    max_ack_pending=1,  # 关键限制参数
    ack_policy=AckPolicy.EXPLICIT
) 

# Nats 
NATS_HOST = '0.0.0.0'
NATS_PORT = '4222'
DOCUMENT_EXT_SUBJECT = "file_process"

MAX_RECONNECT_ATTEMPTS = 10 

## CPU密集型阻塞任务
def time_consuming_task():
    import time; time.sleep(50)
    return

async def task_cb(msg):
    time_consuming_task()  # 阻塞asyncio事件循环
    await msg.ack()
    print("acknowledged document-extraction !")

    
async def run():
    logging.info("Started Document processing Consumer ...")
    
    async def error_cb(e):
        sys.exit()

    async def disconnected_cb():
        logging.info(f"Got disconnected from NATS server .. Retrying .")

    async def reconnected_cb():
        logging.info("Got reconnected...")

    
    nc = await nats.connect(
        servers=f"nats://{NATS_HOST}:{NATS_PORT}",
        error_cb=error_cb,
        reconnected_cb=reconnected_cb,
        disconnected_cb=disconnected_cb,
        max_reconnect_attempts=MAX_RECONNECT_ATTEMPTS,
    )

    # Create JetStream context.
    js = nc.jetstream()
    await js.add_stream(
        subjects=[DOCUMENT_EXT_SUBJECT],
        name="sample-stream", 
        retention=RetentionPolicy.WORK_QUEUE,
    )
 

    await js.subscribe(
        DOCUMENT_EXT_SUBJECT, 
        stream="sample-stream",
        queue = "worker_queue",  # 队列组(分发组)
        cb=task_cb,
        manual_ack=True,
        config=consumer_config,
    )
    
    def signal_handler():
        sys.exit()
    for sig in ('SIGINT', 'SIGTERM'):
        asyncio.get_running_loop().add_signal_handler(getattr(signal, sig), signal_handler)

    await nc.flush()
    logging.info("Done ... ?")


if __name__ == '__main__':
    loop = asyncio.get_event_loop()
    try:
        loop.run_until_complete(run())
        loop.run_forever()
    except Exception as e:
        print("Got error : ", e)
    finally:
        loop.close()

问题描述

启动3个consumer.py进程后,任务并未并行执行:单个进程处理任务时,NATS服务器不会将任务推送给另外两个进程,所有任务逐个依次执行。提前在time_consuming_task()前调用await msg.ack()也无法解决问题。

问题原因及解决方案

1. 同步阻塞任务卡住asyncio事件循环

time_consuming_task()中的time.sleep(50)是同步阻塞操作,会完全占用asyncio的单线程事件循环。这导致:

  • 消费者进程无法处理NATS服务器的新消息推送
  • 无法响应NATS的心跳检测,甚至可能被判定为连接断开
  • 即使提前ack,事件循环被阻塞,消费者也无法接收下一条消息

解决方法:用asyncio.to_thread将阻塞任务放到单独线程执行,避免阻塞事件循环:

async def task_cb(msg):
    # 用线程池执行阻塞任务
    await asyncio.to_thread(time_consuming_task)
    await msg.ack()
    print("acknowledged document-extraction !")

2. max_ack_pending参数限制

当前consumer_config中max_ack_pending=1,意味着每个消费者进程最多只能有1条未确认的消息。即使解决了阻塞问题,单个消费者一次也只能处理1条任务。如果希望每个消费者能并行处理多个任务,可以适当调高该值:

consumer_config = ConsumerConfig(
    ack_wait=900,
    max_deliver=1, 
    max_ack_pending=5,  # 允许每个消费者有5条未确认消息
    ack_policy=AckPolicy.EXPLICIT
) 

3. 队列组的分发逻辑验证

确保所有消费者都使用同一个queue="worker_queue"参数,这样NATS才会将任务均匀分发给队列组内的不同进程。当前代码已正确设置该参数,无需修改,但需要确认所有消费者进程都连接到同一个NATS服务器且配置一致。


内容的提问来源于stack exchange,提问作者bad programmer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 20:14:49