Celery task_success与task_failure信号无法触发问题求助
Celery task_success与task_failure信号无法触发问题求助
Celery Signals such as task_success and task_failure doesn't work.
我有一个部署在Docker中的项目,包含一些异步任务。通过docker-compose运行了redis:alpine、uvicorn和celery三个服务。
Redis的启动命令如下:
docker run redis:alpine -p 6379:6379
Celery的启动命令如下:
python -m celery -A app.infrastructure.celery_tasks worker -E
我的Celery实例配置如下:
# app/infrastructure/celery_tasks/celery_config.py from celery import Celery from kombu import Exchange, Queue app = Celery( "tasks", broker="redis://redis:6379/0", backend="redis://redis:6379/0", ) app.conf.update( broker_connection_retry_on_startup=True, global_retry_backoff=3, CELERY_TASK_ACKS_LATE=True, CELERY_TASK_RETRY_POLICY={ "max_retries": 3, "interval_start": 0, "interval_step": 2, "interval_max": 30, }, use_tz=False, enable_utc=True, worker_heartbeat=3600, broker_transport_options={"visibility_timeout": 3600}, task_serializer="json", accept_content=["json", "application/json"], result_serializer="json", worker_send_task_events=True, task_track_started=True, result_extended=True, task_send_sent_event=True, task_allow_error_cb_on_chord_header=True, task_acks_on_failure_or_timeout=True, ) task_exchange = Exchange("tasks", type="direct") app.conf.task_queues = ( Queue("common", task_exchange, routing_key="tasks.common"), # ..., ) app.conf.task_routes = { "common": {"queue": "common"}, # ..., } app.autodiscover_tasks(["app.infrastructure.celery_tasks"])
__init__.py内容:
# app/infrastructure/celery_tasks/__init__.py from .celery_config import app as celery_app __all__ = ("celery_app",)
我在tasks.py中声明了一个示例任务和task_success信号处理函数(和Celery配置文件分开):
# app/infrastructure/celery_tasks/tasks.py from celery.signals import task_success @shared_task("task_example") def example(user): return {"file": "video.mp4", "user_id": user} @task_success.connect def task_success_handler(sender=None, result=None, **kwargs): print(f"Result: {result}") if result else None
我通过以下代码将任务加入队列:
# main.py from app.infrastructure.celery_tasks.celery_config import app as celery_app task = celery_app.send_task( "task_example", args=("foo"), queue="common", routing_key="tasks.common", ) return {"task_id": task.id}
目前的问题是:没有任何时刻触发task_success事件(信号)。
我猜测可能的问题点在:
- Broker消息传递(Redis)
- Celery配置
- Kombu
有没有什么想法或建议?如果需要更多信息,可以随时问我。
备注:内容来源于stack exchange,提问作者Lucas Zurverra
相关产品推荐
相关产品推荐

