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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 09:04:54