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

如何在Python脚本中启动dramatiq worker并正常处理任务?

问题核心原因

  • 第一是代码执行顺序错误:Worker.start() 是阻塞方法,你将它放在 actor 定义和任务发送代码之前,调用后主线程会直接卡在启动 worker 的步骤,后面的 actor 注册、print_hello_world.send() 代码根本不会执行,worker 自然也感知不到这个 actor 的存在,无法处理对应的任务。
  • 第二是单线程阻塞冲突:如果要在同一个脚本内同时运行 worker 和发送测试任务,需要将 worker 放在独立线程启动,否则主线程被 worker 占满后无法执行后续任务发送逻辑。

修复后的代码

import dramatiq
import threading
import time
from dramatiq.brokers.redis import RedisBroker
from dramatiq.results.backends import RedisBackend
from dramatiq.results import Results
from dramatiq.worker import Worker

# 1. 初始化broker和中间件
redis_broker = RedisBroker(host="127.0.0.1", port=6379)
results_backend = RedisBackend(url="redis://127.0.0.1:6379")
redis_broker.add_middleware(Results(backend=results_backend))
dramatiq.set_broker(redis_broker)

# 2. 先注册actor,必须在worker启动前完成
@dramatiq.actor(queue_name="default", max_retries=1, store_results=True)
def print_hello_world():
    print("Hello World!")

# 3. 初始化worker
worker = Worker(broker=redis_broker)

# 4. 用独立线程启动worker,避免阻塞主线程
worker_thread = threading.Thread(target=worker.start, daemon=True)
worker_thread.start()

# 5. 等待worker完成启动监听,再发送任务
time.sleep(1)

# 6. 发送任务
message = print_hello_world.send()

# 7. 等待任务执行完成,避免脚本直接退出杀死worker
time.sleep(2)

# 可选:停止worker
worker.stop()

验证说明

运行修复后的代码会直接在控制台打印Hello World!,查询Redis可以看到任务执行后的结果缓存key,和用dramatiq test命令启动的运行效果一致。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 08:45:03