Day 19 / 共 20 天 · 阶段6 运维与收官

任务队列与可观测:让慢活跑后台、让每次调用看得见

前 18 天都在讲"处理一次请求"。但生产系统还有两类"看不见却致命"的需求:①有些活很慢(给几百个文档建索引、发一批邮件),不能让用户干等——要甩到后台异步跑;②系统必须"看得见"——一次对话调了几次模型、花了多少 token、哪一步慢,都得能追踪。今天进 api/tasks/(Celery 异步任务)和 core/ops/(可观测追踪),弄清:①Celery 怎么把慢活甩到后台、还按队列分工;②一次对话的追踪数据怎么被攒批、异步上报到 Langfuse 这类平台。

📍 你在 20 天里的位置(阶段6:运维与收官 · D19-20)
D16 工具系统 D17 Agent D18 插件/MCP D19 任务队列/可观测 D20 收官全景
💡 先用两个类比兜住今天 类比一:Celery 任务队列像餐厅的"叫号取餐"。你点了一份要现烤 20 分钟的披萨,收银员不会让你堵在窗口干等——给你个号、你去坐着,做好了叫你。前台(Dify 的 Web 进程)只管"接单、发号"(往队列丢任务),后厨(Celery worker 进程)在后面慢慢烤。用户请求秒回,慢活后台跑。类比二:可观测的"追踪"像快递的物流轨迹。一次对话经过"收件→分拣→运输→派送"好几个环节,追踪系统在每个环节盖个戳(几点到、花了多久),最后拼成一条完整轨迹。出问题时你一看轨迹就知道"卡在哪个环节",而不是对着黑箱瞎猜。
L01

痛点:慢活卡住用户 + 系统是个黑箱

🤔 痛点用户上传一个 500 页的 PDF 建知识库。切片、算 embedding、写向量库,可能要几分钟。如果同步处理,用户的浏览器就得转圈几分钟、甚至超时报错。另一个痛点:线上有人投诉"这次回答特别慢、还很贵",你想查"到底是哪一步慢、哪次调用烧了 token",可日志一片混沌,根本对不上是哪次对话。一个要"慢活异步化",一个要"全链路看得见"。
💡 本质:异步解耦 + 埋点追踪两个成熟套路:①任务队列(Celery + Redis/消息队列)——把慢活封装成"任务"丢进队列,Web 进程立刻返回,独立的 worker 进程在后台消费。请求和执行解耦②可观测(Tracing)——在关键节点埋点,把"这次对话/这次调用"的耗时、token、输入输出等采集下来,攒批异步上报到专业追踪平台(Langfuse、LangSmith 等)。系统从"黑箱"变"玻璃箱"。
L02

Celery:把慢活甩到后台的引擎

Dify 用 Celery 做任务队列,初始化在 extensions/ext_celery.pyextensions/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"],全局共用。还在这里注册了要预加载的任务模块 importsext_celery.py:154)和定时任务 beat_schedule
大白话Celery = "前台 + 后厨 + 叫号系统"三件套。Web 进程是前台(接单发号),worker 进程是后厨(真干活),broker(Redis)是叫号系统(传单子)。三者是独立的,所以你可以单独多开几个后厨(worker)来加速,前台一点不受影响。
L03

@shared_task 与"队列分工"

一个具体任务长这样——看 tasks/document_indexing_task.pytasks/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+ 个任务文件,各归各的队列。
幂等与重试后台任务可能失败重跑,所以这类任务通常设计成"重跑也不出错"(幂等)。这是异步任务的基本素养。
💡 设计取舍:异步的代价是"最终一致"把活甩后台,用户请求快了,但代价是结果不是立刻就绪的——文档"上传成功"不等于"已经能被检索",中间有后台处理的时间差。所以 Dify 的文档有"处理中/已完成"状态,前端要轮询或等通知。用"最终一致"换"响应速度",是所有异步系统都要向用户解释清楚的取舍。
L04

可观测:一次对话怎么被追踪

追踪的总管在 core/ops/ops_trace_manager.py。核心角色两个:TraceTaskcore/ops/ops_trace_manager.py:651,"一条待上报的追踪任务")和 TraceQueueManagerops_trace_manager.py:1491,"攒批调度器")。你在 Day16 见过的 TraceQueueManager、Day18 注入的追踪头,都汇到这里。

埋点在哪模型调用、工具调用(Day16 的小票)、工作流节点、Agent 每一步……关键节点都会产生一个 TraceTask。它记录"这一段发生了什么、耗时多少、token 多少、输入输出是什么"。
为什么不立刻上报如果每产生一条追踪就同步发一次网络请求给 Langfuse,会严重拖慢主流程(一次对话可能产生几十条)。所以要"攒起来、批量、异步"上报——这就是 TraceQueueManager 的活(L05)。
OpsTraceManagerops_trace_manager.py:342)另一个管理类,负责"这个 app 配了哪个追踪平台、拿它的追踪实例"。TraceQueueManager 构造时就靠它 get_ops_trace_instance(app_id)
💡 本质:追踪 = "旁路"采集,绝不阻塞主流程可观测的铁律是"观测本身不能拖垮被观测的系统"。所以 Dify 的追踪是"旁路"的:主流程只管把追踪任务快速塞进一个内存队列就继续跑,真正的"攒批 + 上报"由后台定时器和 Celery 完成。用户对话的速度,完全不受追踪上报快慢的影响。
L05

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:1487TRACE_QUEUE_MANAGER_INTERVAL,默认 5 秒;每批上限 BATCH_SIZE 默认 100),send_to_celeryops_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.Timerops_trace_manager.py:1534start_timer)每 5 秒醒一次,一次捞最多 100 条一起处理——把"几十次小上报"合并成"一次批量",大幅减少网络开销
再转 Celery攒批后还不在当前进程发网络,而是存文件 + 丢 ops_trace 队列,交给 Celery worker 真正上报。两级异步:内存队列削峰、Celery 队列干重活。追踪慢了也绝不影响主服务。
专属队列 ops_trace追踪上报走自己的队列,和 dataset 这种业务队列分开——追踪积压不会拖慢知识库索引,反之亦然。呼应 L03 的队列分工。
追踪数据的两级异步(削峰 + 批量上报) 对话/模型/工具产生 TraceTask 内存队列(put)瞬间返回·不阻塞 定时器每5s攒批一次最多100条 Celery ops_traceworker 真上报 Langfuse / LangSmith / …可视化追踪面板
图注:主流程只塞内存队列(微秒返回)→ 定时器攒批 → 转 Celery 专属队列 → worker 批量推给追踪平台。全程不阻塞对话。
L06

接哪些追踪平台:可插拔的 Provider

Dify 不绑死某一家追踪平台,而是做成"可插拔"。看 OpsTraceProviderConfigMapcore/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 会有开销吗?几乎没有。看 L05 的 add_trace_task:它一开始就判断 if self._enterprise_telemetry_enabled or self.trace_instance——没配追踪平台(trace_instance 为空)就直接跳过,连队列都不塞。所以"埋点代码遍布全身"不代表"没用追踪也有开销"。这又是"不常见情况不给常见情况添负担"的设计(和 Day05 的单 key 快路径同理)。
L07

串起来 + 今日小结

📝 真实值:上传 PDF 建库 + 追踪一次对话 异步侧:用户传 PDF → API 立刻 document_indexing_task.delay(dataset_id, [doc_id]) 返回"处理中" → 任务进 dataset 队列 → worker 后台切片/算 embedding/写向量库(耗时 2 分钟)→ 完成后文档状态变"已完成"。用户全程不卡。追踪侧:另一个用户对话,调了 1 次 LLM + 2 次工具 → 产生 3 个 TraceTaskadd_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)
明日预告 · Day 20(收官):20 天的最后一站!我们把前 19 天的所有模块——请求旅程、模型运行时、工作流、RAG、工具、Agent、插件、异步与可观测——用一张"全景知识地图"串成一条完整的血脉,再顺带看一眼 web/ 前端,最后聊聊"学完之后怎么继续深入 Dify"。有始有终,明天见。
← Day 18 插件/MCP Day 20 · 收官全景 →