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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 11:34:57