Dramatiq异步Actor仅用单Worker,如何利用所有Worker并行执行?
问题场景与解决方案
核心问题
使用Dramatiq并行执行任务时,同步Actor会自动分配到多个Worker线程并行运行,但异步Actor始终在单个Worker线程中串行执行。需要让异步Actor也能利用所有启动的Worker线程,且不想将异步逻辑改为同步后手动创建事件循环。
配置与代码示例
Redis Broker配置
redis_broker = RedisBroker( host=settings.redis_url, port=settings.redis_port, db=settings.redis_db, password=settings.redis_password, middleware=[CurrentMessage()] # 异步场景需补充AsyncIO() ) dramatiq.set_broker(redis_broker)
同步Actor实现
@dramatiq.actor def jobs(param: int): logger.info(f"Started {param}") time.sleep(param) logger.info(f"Ended {param}")
异步Actor实现
@dramatiq.actor async def jobs(param: int): logger.info(f"Started {param}") await asyncio.sleep(param) logger.info(f"Ended {param}")
任务运行方式
g = group([ jobs.send(1), jobs.send(2), jobs.send(3), ]).run()
运行结果对比
同步运行结果(多线程并行)
[2024-10-11 02:38:42,251] [PID 420038] [Thread-3] [dramatiq] [INFO] Started 3 [2024-10-11 02:38:42,251] [PID 420038] [Thread-5] [dramatiq] [INFO] Started 2 [2024-10-11 02:38:42,251] [PID 420038] [Thread-6] [dramatiq] [INFO] Started 3 [2024-10-11 02:38:42,252] [PID 420038] [Thread-7] [dramatiq] [INFO] Started 1 [2024-10-11 02:38:42,252] [PID 420038] [Thread-8] [dramatiq] [INFO] Started 2 [2024-10-11 02:38:42,252] [PID 420038] [Thread-4] [dramatiq] [INFO] Started 1 [2024-10-11 02:38:43,253] [PID 420038] [Thread-7] [dramatiq] [INFO] Ended 1 [2024-10-11 02:38:43,254] [PID 420038] [Thread-4] [dramatiq] [INFO] Ended 1 [2024-10-11 02:38:44,254] [PID 420038] [Thread-5] [dramatiq] [INFO] Ended 2 [2024-10-11 02:38:44,254] [PID 420038] [Thread-8] [dramatiq] [INFO] Ended 2 [2024-10-11 02:38:45,255] [PID 420038] [Thread-6] [dramatiq] [INFO] Ended 3 [2024-10-11 02:38:45,255] [PID 420038] [Thread-3] [dramatiq] [INFO] Ended 3
异步运行结果(单线程串行)
[2024-10-11 02:50:03,245] [PID 422276] [Thread-1] [dramatiq] [INFO] Started 3 [2024-10-11 02:50:03,245] [PID 422276] [Thread-1] [dramatiq] [INFO] Started 2 [2024-10-11 02:50:03,245] [PID 422276] [Thread-1] [dramatiq] [INFO] Started 3 [2024-10-11 02:50:03,245] [PID 422276] [Thread-1] [dramatiq] [INFO] Started 1 [2024-10-11 02:50:03,245] [PID 422276] [Thread-1] [dramatiq] [INFO] Started 2 [2024-10-11 02:50:03,245] [PID 422276] [Thread-1] [dramatiq] [INFO] Started 1 [2024-10-11 02:50:04,246] [PID 422276] [Thread-1] [dramatiq] [INFO] Ended 1 [2024-10-11 02:50:04,246] [PID 422276] [Thread-1] [dramatiq] [INFO] Ended 1 [2024-10-11 02:50:05,247] [PID 422276] [Thread-1] [dramatiq] [INFO] Ended 2 [2024-10-11 02:50:05,247] [PID 422276] [Thread-1] [dramatiq] [INFO] Ended 2 [2024-10-11 02:50:06,246] [PID 422276] [Thread-1] [dramatiq] [INFO] Ended 3 [2024-10-11 02:50:06,246] [PID 422276] [Thread-1] [dramatiq] [INFO] Ended 3
解决方法
1. 启动多线程Worker
运行Dramatiq Worker时显式指定线程数,每个线程会初始化独立的异步事件循环:
dramatiq your_script_module --threads 6
2. 完善Broker中间件配置
异步场景下必须添加AsyncIO中间件,确保异步任务能被正确分发到多线程:
from dramatiq.middleware import AsyncIO, CurrentMessage redis_broker = RedisBroker( host=settings.redis_url, port=settings.redis_port, db=settings.redis_db, password=settings.redis_password, middleware=[CurrentMessage(), AsyncIO()] # 加入AsyncIO中间件 ) dramatiq.set_broker(redis_broker)
3. 处理异步Actor中的阻塞操作
如果异步Actor内需要调用同步阻塞函数,用asyncio.to_thread将其放到线程池执行,避免阻塞事件循环:
@dramatiq.actor async def jobs(param: int): logger.info(f"Started {param}") # 用to_thread包装同步阻塞操作 await asyncio.to_thread(time.sleep, param) logger.info(f"Ended {param}")
调整后,异步Actor任务会被分发到多个Worker线程,每个线程的事件循环并行处理任务,实现和同步Actor一致的多线程并行效果。
内容的提问来源于stack exchange,提问作者sashaaero
相关产品推荐
相关产品推荐

