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
相关产品推荐
相关产品推荐

