Dify 中有不少操作无法在一次 HTTP 请求内快速完成,例如:

  • 解析文件并建立知识库索引;
  • 执行较长的工作流;
  • 安装或升级插件;
  • 批量发送通知;
  • 清理历史数据;
  • 汇总日志和用量。

如果 Web 进程一直等待这些操作,连接容易超时,服务也很快被长任务拖满。因此 Dify 使用 Celery 把一部分工作交给后台 Worker。

但“把函数加上异步装饰器”只是起点。真正可靠的任务系统,还要解决队列隔离、重试一致性、多租户公平、状态可见性和故障恢复。

一、同步请求为什么不适合长任务

同步处理的调用链通常是:

1
客户端 → API → 执行业务 → 返回结果

只要业务没有结束,API 进程和连接就一直被占用。任务耗时从几百毫秒增长到几分钟后,会出现:

  • 网关或浏览器提前超时;
  • Web Worker 被长请求占满;
  • 用户刷新页面后重复提交;
  • 服务重启导致执行中断;
  • 瞬时流量直接压到数据库和模型服务。

引入消息队列后,请求只负责校验、记录状态和投递消息:

flowchart LR
    A[客户端] --> B[API]
    B --> C[(数据库
记录任务)] B --> D[(Broker)] D --> E[Celery Worker] E --> F[解析/索引/工作流] F --> C C --> G[查询状态或推送事件] G --> A

API 可以快速响应,Worker 则按系统承载能力消费任务,流量高峰因此被队列削平。

二、Celery 在其中扮演什么角色

Celery 主要连接三类组件:

  • Producer:API 或定时调度器,负责投递任务;
  • Broker:保存待消费消息,Dify 常用 Redis;
  • Worker:取出消息并执行任务函数。

Dify 还使用 Celery Beat 周期性投递清理、监控、刷新等任务。Beat 只负责“按时发消息”,真正执行工作的仍然是 Worker。

需要特别注意:Broker 中的消息不等于完整业务状态。任务是否成功、文档索引到哪一步,仍应记录在数据库或专门的状态存储中。

三、为什么要拆分多个队列

如果所有任务都进入一个队列,一批耗时很长的文档索引任务可能挡住邮件、日志上报和工作流执行。

Dify 会按用途使用不同队列,当前代码中可以看到 datasetworkflow_storage、工作流执行、API Token 更新等任务通道。拆分队列的价值在于:

  • 不同类型的任务互不阻塞;
  • 可以分别配置 Worker 数量和并发;
  • 可以把 CPU、内存或网络特征不同的任务部署到不同机器;
  • 可以为关键任务设置更高优先级;
  • 故障和积压更容易定位。

一个实用的划分思路是:

队列类型 任务特征 调度重点
知识库索引 CPU、网络和向量写入较重 限制并发,防止挤压存储
工作流执行 时延敏感,可能长时间等待模型 保证专用容量,支持取消
邮件通知 允许短暂延迟 批量处理,限制外部 API 速率
日志与追踪 数量大、单次较轻 批量写入,允许重试
定时清理 不紧急,可能扫描大量数据 低峰执行,严格限制批次

队列不是拆得越多越好。每增加一个队列,都要确保部署中有 Worker 正在监听,否则任务会一直积压却无人处理。

四、任务消息应该携带什么

任务参数应尽量小而明确,通常传递 ID,而不是整个数据库对象:

1
2
3
4
5
6
index_document.delay(
tenant_id=tenant_id,
dataset_id=dataset_id,
document_id=document_id,
operator_id=user_id,
)

Worker 收到任务后重新从数据库读取最新状态。这样可以避免:

  • 序列化复杂对象;
  • 消息中残留过期数据;
  • 把密钥或大段文本写入 Broker;
  • 代码升级后旧对象无法反序列化。

同时,租户 ID 或能够可靠推导租户的资源 ID 不能缺失。后台进程没有 HTTP 请求上下文,必须自己重建租户边界。

五、异步任务必须允许重复执行

消息系统很难保证任务“绝对只执行一次”。Worker 可能已经完成业务写入,却在确认消息前崩溃,Broker 随后会再次投递。

因此任务最好具备幂等性:相同任务执行多次,最终结果保持一致。

常见做法包括:

  • 为业务操作设置唯一任务 ID;
  • 执行前检查目标状态是否已经完成;
  • 数据库使用唯一约束防止重复插入;
  • 更新时采用条件语句或版本号;
  • 外部请求携带幂等键;
  • 对通知类副作用单独记录发送状态。

例如“重建文档索引”可以先检查文档是否仍存在、是否已经被新版本任务处理,再决定继续还是直接结束。

六、重试不是捕获异常后再跑一次

并非所有失败都应该重试:

  • 网络抖动、限流和临时超时通常可以重试;
  • 参数错误、权限不足和文件格式错误通常不能靠重试恢复;
  • 数据库约束异常要先判断是否已成功执行过;
  • 模型额度耗尽可能需要等待人工处理。

合理的重试策略一般包含:

1
2
3
4
5
可重试异常白名单
+ 最大次数
+ 指数退避
+ 随机抖动
+ 最终失败状态

如果所有异常都立即重试,外部服务故障时会形成“重试风暴”,进一步压垮系统。

Dify 的部分任务会显式配置 max_retries 和延迟;另一些任务则通过数据库状态和业务级重试入口恢复。两种方式都需要明确最终失败后由谁发现、谁处理。

七、多租户场景为什么需要二级排队

Celery 队列解决不同任务类型之间的调度,但同一 dataset 队列中,仍可能有一个租户一次提交成千上万个任务。

Dify 的 RAG 流程还提供租户隔离任务队列:使用 Redis List 保存某个租户的等待任务,并设置租户级并发限制。可以把它理解为:

1
2
Celery 队列:控制全局由哪些 Worker 执行
租户队列:控制同一租户一次可以占用多少执行名额
flowchart TD
    A[租户 A 的索引请求] --> A1[租户 A 等待队列]
    B[租户 B 的索引请求] --> B1[租户 B 等待队列]
    C[租户 C 的索引请求] --> C1[租户 C 等待队列]
    A1 --> D[全局 dataset 队列]
    B1 --> D
    C1 --> D
    D --> E1[Worker 1]
    D --> E2[Worker 2]
    D --> E3[Worker 3]

这样既能保护系统,也能避免大客户完全饿死小客户的任务。

八、任务状态应该如何设计

用户需要知道任务是“等待中”“执行中”还是“失败了”,因此业务表通常要有明确状态:

1
2
3
PENDING → RUNNING → COMPLETED
↘ FAILED
↘ CANCELLED

状态更新需要考虑异常退出:如果 Worker 在写入 RUNNING 后突然断电,任务可能永远停在执行中。

可以配合以下机制恢复:

  • 保存开始时间、心跳和最后更新时间;
  • 定时扫描超时任务;
  • 区分“可重试失败”和“永久失败”;
  • 记录错误摘要和重试次数;
  • 为人工重试提供清晰入口;
  • 定期清理过期状态和孤儿数据。

进度条也不应凭空估算。能按文档数、分段数计算时才显示百分比,无法估计时使用阶段状态更诚实。

九、数据库事务与投递消息的缝隙

一个经典问题是:

  1. 数据库创建任务记录;
  2. 提交事务;
  3. 向 Celery 投递消息。

如果第 2 步成功而第 3 步失败,数据库里会留下永远没人执行的任务。反过来,如果先投递消息,Worker 又可能在数据库事务提交前读取不到记录。

常见解决方案包括:

  • 事务提交成功后再投递,并对投递失败进行补偿;
  • 使用 Outbox 表,在同一事务内记录待发布事件;
  • 由独立进程扫描 Outbox 并可靠投递;
  • 定时扫描长时间处于 PENDING 的孤儿任务。

任务量不大时可以从“提交后投递+定时补偿”起步,关键是承认这个缝隙存在,而不是假设两套系统能天然原子提交。

十、如何观察任务系统是否健康

只看 Worker 进程是否存活远远不够。建议监控:

  • 各队列等待长度;
  • 最老消息等待时间;
  • 任务吞吐量和执行时长;
  • 成功率、重试率和永久失败数;
  • 各租户积压量;
  • Worker 并发、CPU、内存和数据库连接;
  • Broker 容量、连接数和延迟。

Dify 提供知识库队列监控阈值和周期配置,可在积压超过阈值时告警。相比单纯的队列长度,“最老任务等了多久”通常更能反映用户体验。

日志中至少应带上:

1
task_id、task_name、tenant_id、resource_id、retry_count

但不要把文档原文、模型密钥或完整任务参数直接写入日志。

十一、扩容与优雅停机

队列积压时,第一反应往往是增加 Worker。扩容前还要确认瓶颈在哪里:

  • CPU 密集任务适合增加进程或机器;
  • 大量外部 I/O 可适度增加并发;
  • 数据库连接已满时,继续加 Worker 只会恶化故障;
  • 模型供应商限流时,需要控制消费速度而不是盲目扩容。

发布或缩容时也不能直接杀死 Worker。优雅停机应停止领取新任务,等待正在执行的任务完成或安全中断,再退出进程。对长工作流,还需要专门的取消、续跑或检查点机制。

十二、总结

Dify 的异步任务系统可以理解为三层治理:

1
2
3
Celery 与 Broker:解耦请求和执行
多队列与专用 Worker:隔离不同任务类型
租户队列与业务状态:保证公平和可恢复

稳定的异步系统不是“任务最终大概会跑完”,而是能够回答:任务在哪里、为什么等待、能否安全重试、失败后如何恢复,以及任何一个租户是否正在挤占整个系统。