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 会按用途使用不同队列,当前代码中可以看到 dataset、workflow_storage、工作流执行、API Token 更新等任务通道。拆分队列的价值在于:
- 不同类型的任务互不阻塞;
- 可以分别配置 Worker 数量和并发;
- 可以把 CPU、内存或网络特征不同的任务部署到不同机器;
- 可以为关键任务设置更高优先级;
- 故障和积压更容易定位。
一个实用的划分思路是:
| 队列类型 | 任务特征 | 调度重点 |
|---|---|---|
| 知识库索引 | CPU、网络和向量写入较重 | 限制并发,防止挤压存储 |
| 工作流执行 | 时延敏感,可能长时间等待模型 | 保证专用容量,支持取消 |
| 邮件通知 | 允许短暂延迟 | 批量处理,限制外部 API 速率 |
| 日志与追踪 | 数量大、单次较轻 | 批量写入,允许重试 |
| 定时清理 | 不紧急,可能扫描大量数据 | 低峰执行,严格限制批次 |
队列不是拆得越多越好。每增加一个队列,都要确保部署中有 Worker 正在监听,否则任务会一直积压却无人处理。
四、任务消息应该携带什么
任务参数应尽量小而明确,通常传递 ID,而不是整个数据库对象:
1 | index_document.delay( |
Worker 收到任务后重新从数据库读取最新状态。这样可以避免:
- 序列化复杂对象;
- 消息中残留过期数据;
- 把密钥或大段文本写入 Broker;
- 代码升级后旧对象无法反序列化。
同时,租户 ID 或能够可靠推导租户的资源 ID 不能缺失。后台进程没有 HTTP 请求上下文,必须自己重建租户边界。
五、异步任务必须允许重复执行
消息系统很难保证任务“绝对只执行一次”。Worker 可能已经完成业务写入,却在确认消息前崩溃,Broker 随后会再次投递。
因此任务最好具备幂等性:相同任务执行多次,最终结果保持一致。
常见做法包括:
- 为业务操作设置唯一任务 ID;
- 执行前检查目标状态是否已经完成;
- 数据库使用唯一约束防止重复插入;
- 更新时采用条件语句或版本号;
- 外部请求携带幂等键;
- 对通知类副作用单独记录发送状态。
例如“重建文档索引”可以先检查文档是否仍存在、是否已经被新版本任务处理,再决定继续还是直接结束。
六、重试不是捕获异常后再跑一次
并非所有失败都应该重试:
- 网络抖动、限流和临时超时通常可以重试;
- 参数错误、权限不足和文件格式错误通常不能靠重试恢复;
- 数据库约束异常要先判断是否已成功执行过;
- 模型额度耗尽可能需要等待人工处理。
合理的重试策略一般包含:
1 | 可重试异常白名单 |
如果所有异常都立即重试,外部服务故障时会形成“重试风暴”,进一步压垮系统。
Dify 的部分任务会显式配置 max_retries 和延迟;另一些任务则通过数据库状态和业务级重试入口恢复。两种方式都需要明确最终失败后由谁发现、谁处理。
七、多租户场景为什么需要二级排队
Celery 队列解决不同任务类型之间的调度,但同一 dataset 队列中,仍可能有一个租户一次提交成千上万个任务。
Dify 的 RAG 流程还提供租户隔离任务队列:使用 Redis List 保存某个租户的等待任务,并设置租户级并发限制。可以把它理解为:
1 | 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 | PENDING → RUNNING → COMPLETED |
状态更新需要考虑异常退出:如果 Worker 在写入 RUNNING 后突然断电,任务可能永远停在执行中。
可以配合以下机制恢复:
- 保存开始时间、心跳和最后更新时间;
- 定时扫描超时任务;
- 区分“可重试失败”和“永久失败”;
- 记录错误摘要和重试次数;
- 为人工重试提供清晰入口;
- 定期清理过期状态和孤儿数据。
进度条也不应凭空估算。能按文档数、分段数计算时才显示百分比,无法估计时使用阶段状态更诚实。
九、数据库事务与投递消息的缝隙
一个经典问题是:
- 数据库创建任务记录;
- 提交事务;
- 向 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 | Celery 与 Broker:解耦请求和执行 |
稳定的异步系统不是“任务最终大概会跑完”,而是能够回答:任务在哪里、为什么等待、能否安全重试、失败后如何恢复,以及任何一个租户是否正在挤占整个系统。