try-except在任务函数内直接捕获异常可防止崩溃、保留上下文并便于排查;需包裹关键逻辑、记录完整错误信息,再根据业务决定重试或软失败。

任务函数内部直接 try-except 捕获异常
大多数情况下,异常发生在任务函数执行过程中,比如数据库查询失败、HTTP 请求超时、JSON 解析出错。这时候最直接有效的方式是在 task 函数体内用 try...except 包裹关键逻辑,避免异常向上抛出导致任务状态变为 FAILURE 而丢失上下文。
常见错误现象:任务日志只显示 Task xyz raised unexpected: ValueError(...),但没打印具体参数或中间状态,排查困难。
- 把关键业务逻辑(如
requests.get()、json.loads()、model.save())放在try块里 - 在
except中记录完整错误信息(含traceback.format_exc())和输入参数 - 根据业务决定是否重试(调用
self.retry())或转为“软失败”(返回None或标记字段)
@app.task(bind=True, autoretry_for=(requests.RequestException,), retry_kwargs={'max_retries': 3})
def fetch_user_data(self, user_id):
try:
resp = requests.get(f"https://api.example.com/users/{user_id}", timeout=5)
resp.raise_for_status()
return resp.json()
except Exception as e:
logger.error(f"fetch_user_data failed for {user_id}: {e}", exc_info=True)
raise # 让 Celery 记录 FAILURE 状态,便于监控告警
通过 task_failure 信号监听所有失败任务
当需要统一处理所有任务的异常(比如发告警、写审计日志、触发补偿流程),不能依赖每个任务手动捕获,而应使用 Celery 的 task_failure 信号。它会在任务明确失败(即进入 FAILURE 状态)后触发,且携带完整的异常信息。
注意:该信号不会捕获被 try/except 吞掉且未重新抛出的异常,也不会响应 REJECTED 或 REVOKED 状态。
- 信号处理器必须在 Celery app 初始化后、worker 启动前注册(通常放在 tasks.py 底部或 signals.py 中)
-
exception参数是原始异常对象,可直接用repr(e)或traceback.format_exception()格式化 -
args和kwargs是调用任务时传入的参数,可用于定位问题实例
from celery import signals
<p>@signals.task_failure.connect
def on_task_failure(sender, task_id, exception, args, kwargs, traceback, einfo, **extras):
logger.error(
f"Task {sender.name}[{task_id}] failed: {exception}",
extra={"task_args": args, "task_kwargs": kwargs, "traceback": traceback}
)
if "timeout" in str(exception).lower():
alert_timeout_task(task_id, args)
</p>
使用 Task.on_failure() 方法定制单个任务的失败行为
如果某个任务需要区别于全局策略的失败处理(例如不告警、只写本地文件、或调用特定回调),可以继承 Task 并重写 on_failure() 方法。这个方法在任务抛出异常且未被 retry() 拦截时调用,比信号更轻量、更靠近任务本身。
容易踩的坑:on_failure() 不会自动记录日志,也不影响任务状态(状态仍是 FAILURE),需自行确保可观测性;且它不接收 args/kwargs 的解包形式,需从 self.request 中提取。
-
self.request.args和self.request.kwargs可获取原始参数 -
exc是异常对象,traceback是字符串格式的栈跟踪 - 不要在
on_failure()中调用阻塞操作(如同步 HTTP 请求),可能拖慢 worker
class ReportTask(app.Task):
def on_failure(self, exc, task_id, args, kwargs, einfo):
# args/kwargs 是原始传入值,einfo.traceback 是字符串栈
with open(f"/tmp/failures/{task_id}.log", "w") as f:
f.write(f"Failed at {datetime.now()}\n{einfo.traceback}")
<p>@app.task(base=ReportTask)
def generate_monthly_report(month):
...
</p>
异步链式任务中异常传播与断点控制
用 chord、group 或 | 链接多个任务时,上游异常默认不会中断整个链(除非显式配置),下游仍可能被执行,造成数据不一致或重复动作。
典型问题:一个 chord 的 header 里某个任务失败,callback 却仍被调用,且收不到失败任务的 ID 和错误信息。
- 对关键链路,用
ignore_result=False(默认)并检查result.get(propagate=False)返回的state字段 - 在 callback 中遍历
result.parent.results判断各子任务状态,跳过FAILED项 - 避免在
chordcallback 里直接调用.get(),否则会阻塞并可能掩盖上游异常类型
真正难处理的是“部分失败+需人工介入”的场景——Celery 本身不提供事务回滚,得靠幂等设计 + 外部状态机来兜底。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











