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

