任务队列与可观测:让慢活跑后台、让每次调用看得见
前 18 天都在讲"处理一次请求"。但生产系统还有两类"看不见却致命"的需求:①有些活很慢(给几百个文档建索引、发一批邮件),不能让用户干等——要甩到后台异步跑;②系统必须"看得见"——一次对话调了几次模型、花了多少 token、哪一步慢,都得能追踪。今天进 api/tasks/(Celery 异步任务)和 core/ops/(可观测追踪),弄清:①Celery 怎么把慢活甩到后台、还按队列分工;②一次对话的追踪数据怎么被攒批、异步上报到 Langfuse 这类平台。
痛点:慢活卡住用户 + 系统是个黑箱
Celery:把慢活甩到后台的引擎
Dify 用 Celery 做任务队列,初始化在 extensions/ext_celery.py(extensions/ext_celery.py:99):
# api/extensions/ext_celery.py:99
def init_app(app: DifyApp) -> Celery:
class FlaskTask(Task): # ★ 让每个后台任务也能用 Flask 上下文
def __call__(self, *args, **kwargs):
from core.logging.context import init_request_context
with app.app_context(): # 进入应用上下文(能访问 db 等)
init_request_context()
return self.run(*args, **kwargs)
celery_app = Celery(
app.name, task_cls=FlaskTask,
broker=dify_config.CELERY_BROKER_URL, # ★ broker:任务往哪丢(如 Redis)
backend=dify_config.CELERY_BACKEND, # backend:结果存哪
)
celery_app.conf.update(
...
task_ignore_result=True, # 多数任务不关心返回值,省存储
task_annotations=dify_config.CELERY_TASK_ANNOTATIONS,
)
broker"消息中间人"。Web 进程把任务丢给 broker(Dify 默认用 Redis),worker 从 broker 取任务执行。broker 就是那个"叫号系统":前台把号丢进去,后厨来取。FlaskTaskCelery worker 是独立进程,默认没有 Flask 的应用上下文。这个包装让每个任务运行前先 with app.app_context()——这样后台任务里也能正常用数据库、配置。不然任务一碰 db 就报"working outside of application context"。task_ignore_result=True大多数任务(建索引、发邮件)你不关心它"返回了什么",只关心"做没做"。忽略结果就不用把结果写进 backend,省资源。init_app 返回 celery_app这个实例被挂到 app.extensions["celery"],全局共用。还在这里注册了要预加载的任务模块 imports(ext_celery.py:154)和定时任务 beat_schedule。@shared_task 与"队列分工"
一个具体任务长这样——看 tasks/document_indexing_task.py(tasks/document_indexing_task.py:32):
# api/tasks/document_indexing_task.py:32
@shared_task(queue="dataset") # ★ 声明这是个 Celery 任务,丢进 "dataset" 队列
def document_indexing_task(dataset_id: str, document_ids: list):
"""
Async process document
Usage: document_indexing_task.delay(dataset_id, document_ids) # ★ 调用方用 .delay() 异步触发
"""
... # 切片、建索引的慢活都在这里,后台跑
@shared_task装饰器。加上它,一个普通函数就变成"可被丢进队列后台执行"的任务。.delay(...)★关键动作。业务代码不直接调 document_indexing_task(...)(那是同步、会阻塞),而是调 document_indexing_task.delay(...)——这行瞬间返回,任务被丢进队列,真正执行在 worker。用户请求秒回。queue="dataset"★队列分工。Dify 把任务按类型分到不同队列:dataset(知识库索引)、ops_trace(追踪上报)、邮件、清理等。这样可以给"重活队列"多配 worker、给"轻活队列"少配,互不抢资源。api/tasks/ 下有 50+ 个任务文件,各归各的队列。幂等与重试后台任务可能失败重跑,所以这类任务通常设计成"重跑也不出错"(幂等)。这是异步任务的基本素养。可观测:一次对话怎么被追踪
追踪的总管在 core/ops/ops_trace_manager.py。核心角色两个:TraceTask(core/ops/ops_trace_manager.py:651,"一条待上报的追踪任务")和 TraceQueueManager(ops_trace_manager.py:1491,"攒批调度器")。你在 Day16 见过的 TraceQueueManager、Day18 注入的追踪头,都汇到这里。
埋点在哪模型调用、工具调用(Day16 的小票)、工作流节点、Agent 每一步……关键节点都会产生一个 TraceTask。它记录"这一段发生了什么、耗时多少、token 多少、输入输出是什么"。为什么不立刻上报如果每产生一条追踪就同步发一次网络请求给 Langfuse,会严重拖慢主流程(一次对话可能产生几十条)。所以要"攒起来、批量、异步"上报——这就是 TraceQueueManager 的活(L05)。OpsTraceManager(ops_trace_manager.py:342)另一个管理类,负责"这个 app 配了哪个追踪平台、拿它的追踪实例"。TraceQueueManager 构造时就靠它 get_ops_trace_instance(app_id)。TraceQueueManager:攒批 + 异步上报
看 TraceQueueManager 怎么"攒"和"发"(core/ops/ops_trace_manager.py:1506):
# api/core/ops/ops_trace_manager.py:1506
def add_trace_task(self, trace_task: TraceTask):
global trace_manager_timer, trace_manager_queue
try:
if self._enterprise_telemetry_enabled or self.trace_instance:
trace_task.app_id = self.app_id
trace_manager_queue.put(trace_task) # ★① 只是塞进内存队列,瞬间返回(不阻塞对话)
except Exception:
logger.exception("Error adding trace task, trace_type %s", trace_task.trace_type)
finally:
self.start_timer() # ② 确保后台定时器在跑
def collect_tasks(self): # ③ 定时器触发时:从队列捞一批
tasks = []
while len(tasks) < trace_manager_batch_size and not trace_manager_queue.empty():
task = trace_manager_queue.get_nowait()
tasks.append(task)
return tasks
def run(self):
tasks = self.collect_tasks()
if tasks:
self.send_to_celery(tasks) # ★④ 一批一起交给 Celery 异步上报
定时器每隔几秒跑一次(间隔来自 ops_trace_manager.py:1487 的 TRACE_QUEUE_MANAGER_INTERVAL,默认 5 秒;每批上限 BATCH_SIZE 默认 100),send_to_celery(ops_trace_manager.py:1542)把追踪数据存成文件、再丢 Celery 的 ops_trace 队列:
# api/tasks/ops_trace_task.py:41
@shared_task(queue="ops_trace") # ★ 追踪上报有专属队列,和业务任务隔离
def process_trace_tasks(self, file_info):
... # worker 里真正把数据推给 Langfuse/LangSmith
put 进内存队列主流程调 add_trace_task 时,只做一件事:把任务丢进进程内的队列。微秒级返回,对话流程零感知。定时器 + 攒批后台一个 threading.Timer(ops_trace_manager.py:1534 的 start_timer)每 5 秒醒一次,一次捞最多 100 条一起处理——把"几十次小上报"合并成"一次批量",大幅减少网络开销。再转 Celery攒批后还不在当前进程发网络,而是存文件 + 丢 ops_trace 队列,交给 Celery worker 真正上报。两级异步:内存队列削峰、Celery 队列干重活。追踪慢了也绝不影响主服务。专属队列 ops_trace追踪上报走自己的队列,和 dataset 这种业务队列分开——追踪积压不会拖慢知识库索引,反之亦然。呼应 L03 的队列分工。接哪些追踪平台:可插拔的 Provider
Dify 不绑死某一家追踪平台,而是做成"可插拔"。看 OpsTraceProviderConfigMap(core/ops/ops_trace_manager.py:223 附近),按需懒加载各家的配置和追踪实现:
# api/core/ops/ops_trace_manager.py (Provider 映射节选)
match key:
case TracingProviderEnum.LANGFUSE: # ★ Langfuse
from dify_trace_langfuse.langfuse_trace import LangFuseDataTrace
return {"config_class": LangfuseConfig,
"secret_keys": ["public_key", "secret_key"],
"other_keys": ["host", "project_key"],
"trace_instance": LangFuseDataTrace}
case TracingProviderEnum.LANGSMITH: # ★ LangSmith
from dify_trace_langsmith.langsmith_trace import LangSmithDataTrace
return {"config_class": LangSmithConfig,
"secret_keys": ["api_key"],
"other_keys": ["project", "endpoint"],
"trace_instance": LangSmithDataTrace}
case TracingProviderEnum.OPIK: ... # 还有 Opik 等
按 key match你在 app 里配了哪家追踪平台,就用哪个分支——每家有自己的配置类、密钥字段、追踪实现类。加一家新平台,只要加一个 case。这是"策略模式 + 可插拔"。懒加载 import每个 case 里才 import 对应的库。你没用 LangSmith,就不会加载它的依赖——省内存、也避免"没装某库就启动失败"。secret_keys / other_keys区分"敏感字段(要加密存)"和"普通字段"。密钥类信息(如 secret_key)会被加密处理——安全意识再次出现。add_trace_task:它一开始就判断 if self._enterprise_telemetry_enabled or self.trace_instance——没配追踪平台(trace_instance 为空)就直接跳过,连队列都不塞。所以"埋点代码遍布全身"不代表"没用追踪也有开销"。这又是"不常见情况不给常见情况添负担"的设计(和 Day05 的单 key 快路径同理)。串起来 + 今日小结
document_indexing_task.delay(dataset_id, [doc_id]) 返回"处理中" → 任务进 dataset 队列 → worker 后台切片/算 embedding/写向量库(耗时 2 分钟)→ 完成后文档状态变"已完成"。用户全程不卡。追踪侧:另一个用户对话,调了 1 次 LLM + 2 次工具 → 产生 3 个 TraceTask → add_trace_task 塞进内存队列(对话零延迟)→ 5 秒后定时器攒批 → 丢 ops_trace 队列 → worker 推给 Langfuse → 你在 Langfuse 面板看到这次对话耗时 3.2s、用了 1850 token、哪一步最慢,一目了然。👶 小白:为什么追踪要"内存队列 + Celery"两层异步,一层不够吗?
👨🏫 老师:分工不同。内存队列是"削峰"——对话瞬间产生一堆追踪,先在进程内攒着、微秒级返回,绝不拖慢对话。Celery 队列是"干重活并跨进程"——真正发网络给 Langfuse 可能慢、可能失败要重试,这种活交给专门的 worker 进程最合适,还能和内存队列所在的 Web 进程解耦。一层负责"快", 一层负责"稳", 合起来才既不阻塞又可靠。
🧠 今天你应该能回答
- Celery 三件套是什么?(前台 Web / 后厨 worker / 叫号 broker)
.delay()做了什么?(把任务丢队列、瞬间返回,真正执行在 worker)- 为什么要
queue="dataset"这种队列分工?(重活轻活分开配 worker,互不抢资源) FlaskTask为什么必要?(让后台任务也有 Flask 上下文,能用 db)- 追踪为什么要两级异步?(内存队列削峰不阻塞 + Celery 队列干重活跨进程)
- 没配追踪平台会有开销吗?(几乎没有,
add_trace_task直接跳过)
✋ 10 分钟动手
cd /Users/bitmart/work/codes/github/AI_WORK/dify
# 1. Celery 初始化
sed -n '99,160p' api/extensions/ext_celery.py # init_app / FlaskTask / imports
# 2. 一个真实任务
sed -n '32,45p' api/tasks/document_indexing_task.py # @shared_task(queue="dataset")
ls api/tasks/ | head -30 # 50+ 个后台任务
# 3. 可观测攒批上报
sed -n '1491,1545p' api/core/ops/ops_trace_manager.py # TraceQueueManager
sed -n '41,48p' api/tasks/ops_trace_task.py # process_trace_tasks (queue=ops_trace)
web/ 前端,最后聊聊"学完之后怎么继续深入 Dify"。有始有终,明天见。