数据截至 (上游 commit 5053c08115bd)
运行时:Celery 工蜂群、Redis 协调、多租户与开源分层
30 秒导读: 前五章讲的是「一次请求内部发生了什么」。本章讲这些代码到底跑在哪个进程里、谁按点叫醒它们、挂了之后谁来收尸——一个由 supervisord 托管的 Celery 工蜂群、一套 Redis 锁与栅栏、一层 Postgres schema 级的租户隔离,以及 MIT 主干与企业版之间那个只有一个函数宽的接缝。
1. 这是什么(零基础也能懂)
一句话定义: Onyx 的运行时 = 三类长驻进程(API 服务器、后台工蜂群、模型推理服务器)+ 两个中枢(Postgres 存事实,Redis 做协调)。
前五章里那些漂亮的东西——对话循环、上下文装配、工具与子 agent、检索与引用、索引管线——都不是凭空运行的。它们分居在两种进程里:
| 前面章节的能力 | 实际跑在哪 |
|---|---|
| 一次对话(01/02/03/04) | api_server 进程内的一个 HTTP 流式请求 |
| 索引管线(05) | docfetching / docprocessing 两种 Celery worker |
| 嵌入模型、reranker 推理 | 独立的 model_server 进程 |
| 权限同步、剪枝、清理、监控 | light / heavy / monitoring 等 worker |
它要解决的问题: 用户按下回车要毫秒级响应,而爬 Confluence 一整个空间要几小时。这两类工作不能挤在同一个进程、同一个池子里,否则一次全量索引会把聊天卡死。所以 Onyx 把慢活全部推到消息队列后面,并按「快 / 慢 / 抓 / 算」分成不同的 worker 池,各配各的并发。
一句话直觉: 把 Onyx 想成一个蜂巢——api_server 是接待窗口,Redis 是墙上那块任务板 + 占位牌,Postgres 是账本,supervisord 下的一群 Celery worker 是不同工种的工蜂,Beat 是每隔十几秒敲一次钟的钟楼。
2. 顶层全景(它大概怎么转)
怎么读这张图: 从左到右是"人的请求",从上到下是"机器自己给自己派的活"。中间那条竖线是 Redis——所有跨进程的协调都从这里过。
浏览器 / Slack / API
│
▼
┌──────────────────┐ ┌──────────────────┐
│ api_server │───────▶│ model_server │ 嵌入 / rerank 推理
│ (FastAPI 单进程)│ │ inference/index │ 两份,独立扩缩
└────────┬─────────┘ └──────────────────┘
│ send_task
═════════▼═══════════════════════════════════════ Redis
│ ① broker 队列 ② 锁与 fence ③ 缓存/会话
═════════╤═══════════════════════════════════════
│
┌────────┴──────────────────────────────────────┐
│ background 容器(supervisord 托管) │
│ │
│ 钟楼 Beat ──┬─▶ primary ──派活──┐ │
│ │ (单例协调者) │ │
│ watchdog ───┘ ▼ │
│ light / heavy / docfetching │
│ docprocessing / monitoring │
│ user_file_processing / … │
└───────────────────────┬───────────────────────┘
▼
Postgres(事实与账本)
+ OpenSearch/Vespa(向量索引)
部件职责一览:
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
api_server | FastAPI 应用,挂 60+ 个 router、认证栈、中间件 | backend/onyx/main.py:492 get_application |
| Beat(钟楼) | 按节律往队列里丢"检查"任务,多租户时逐租户展开 | backend/onyx/background/celery/apps/beat.py:26 DynamicTenantScheduler |
| primary(协调者) | 唯一消费默认 celery 队列的 worker,扫描状态并派活 | backend/onyx/background/celery/apps/primary.py:52 |
| 各类 worker | 真正干活;按队列分工,各自并发配置 | backend/onyx/background/celery/apps/*.py |
| watchdog | 盯 Redis 心跳键,Beat 死了就重启它 | backend/onyx/utils/supervisord_watchdog.py:17 main |
model_server | 只做模型推理的 HTTP 服务,与业务代码零耦合 | backend/model_server/main.py:118 get_model_app |
主线走一遍(不进代码): Beat 每 15 秒喊一次 check_for_indexing → primary 接住 ,扫 Postgres 找出该索引的连接器 → 给每个连接器创建一条 IndexAttempt 记录(这就是"占位")并往 connector_doc_fetching 队列丢任务 → docfetching worker 抓文档、按批存进文件存储、每批再往 docprocessing 队列丢一个任务 → docprocessing worker 做切块、嵌入、写索引。
3. Celery 工蜂群:worker → queue 的映射
3.1 分工表
这是本章最该背下来的一张表。队列名来自 backend/onyx/configs/constants.py:442 的 OnyxCeleryQueues,-Q 绑定来自 backend/supervisord.conf,并发默认值来自 backend/onyx/configs/app_configs.py。
| worker | 消费的队列(-Q) | 默认并发 | prefetch | 干什么 |
|---|---|---|---|---|
primary | celery | 4 | 1 | 单例协调者:扫状态、派活、清 Redis |
light | vespa_metadata_sync, connector_deletion, doc_permissions_upsert, checkpoint_cleanup, index_attempt_cleanup, opensearch_migration | 24 | 8 | 短平快:元数据同步、删文档、清检查点 |
heavy | connector_pruning, connector_doc_permissions_sync, connector_external_group_sync, csv_generation, sandbox | 4 | 1 | 长耗时、打外部 API:剪枝、权限同步、CSV 导出、沙箱 |
docfetching | connector_doc_fetching | 1 | 1 | 跑连接器拉原始文档,按批落盘 |
docprocessing | docprocessing | 6 | 1 | 取批次跑索引管线(切块、嵌入、入库) |
user_file_processing | user_file_processing, user_file_project_sync, user_file_delete | 2 | 1 | 用户上传的文件 |
scheduled_tasks | scheduled_tasks | 4 | 1 | Craft 定时 agent 执行器(长跑 LLM 调用) |
monitoring | monitoring | 1 | 1 | 队列深度、连接器成败、进程内存等指标 |
并发默认值锚点:app_configs.py:897(light=24)、:908(light prefetch=8)、:921(docprocessing=6)、:935(docfetching=1)、:949(primary=4)、:958(heavy=4)、:962(monitoring=1)、:966(user_file=2)、:974(scheduled_tasks=4)。
light 是唯一开大 prefetch 的:24 并发 × 8 预取 = 同时在飞 192 个任务(backend/onyx/background/celery/configs/light.py:25-27)。其余全是 prefetch_multiplier = 1——宁可空转也不让一个线程囤活。
3.2 为什么全用线程池,不用进程池
每个 configs/*.py 都写死 worker_pool = "threads"。原因写在 backend/supervisord.conf:20-29 的注释里:Celery + SQLAlchemy 在 prefork 池下会间歇性 WorkerLostError: signal 11 (SIGSEGV)(上游 issue celery#7007)。代价是 worker 吃不到多核;Onyx 的判断是"任务大多在等 Vespa / Postgres 的 I/O",所以能接受。
注意这里有个例外:docfetching 并不真的在线程里跑连接器。docfetching_proxy_task 会另起一个进程执行 docfetching_task,跑完直接 os._exit(0)(backend/onyx/background/celery/tasks/docfetching/tasks.py:332,进程末尾 os._exit(0) 在 :278)。连接器是最容易泄漏内存和吞 CPU 的一环,所以给它单独一条命。
3.3 versioned_apps 这层壳
supervisord 启的不是 apps/primary,而是 versioned_apps/primary。这个模块只有 12 行:
# backend/onyx/background/celery/versioned_apps/primary.py:8-12(真实源码,节选)
set_is_ee_based_on_env_variable()
app: Celery = fetch_versioned_implementation(
"onyx.background.celery.apps.primary",
"celery_app",
)
先按环境变量决定"我是不是企业版",再决定去 onyx. 还是 ee.onyx. 下取那个 celery_app。同一条启动命令,能拉起 MIT 版或 EE 版——这就是第 8 节要展开的可插拔机制在进程入口处的第一次出现。
3.4 README 与代码的三处不一致
backend/onyx/background/README.md 的 worker→queue 表按此 commit 核对下来有三处对不上,读的时候要以代码为准:
| README 说 | 代码实际 | 依据 |
|---|---|---|
有个 Background (consolidated) worker,文件 apps/background.py | 该文件不存在,apps/ 下只有 12 个模块 | ls backend/onyx/background/celery/apps/ |
light 消费 5 个队列 | supervisord 里是 7 个,多 opensearch_migration 和 chat_ttl_deletion | backend/supervisord.conf:46-47 |
未提及 scheduled_tasks worker | 存在且有独立 app/config/supervisord program | backend/supervisord.conf:91-100 |
4. Beat:整个系统的节律
4.1 定时表
Beat 自己不干活,它只按点往队列里丢"去检查一下"的任务。自托管模式下的完整节律(backend/onyx/background/celery/tasks/beat_schedule.py:42-225):
| 任务 | 周期 | 落到哪个队列 | 派生出什么 |
|---|---|---|---|
check_for_indexing | 15s | celery(primary) | → connector_doc_fetching |
check_for_vespa_sync_task | 20s | celery | → vespa_metadata_sync |
check_for_pruning | 20s | celery | → connector_pruning |
check_for_connector_deletion | 20s | celery | → connector_deletion |
check_for_user_file_processing / _project_sync / _delete | 各 20s | celery | → user_file_* |
dispatch_due_scheduled_tasks | 30s | celery | → scheduled_tasks |
check_for_index_attempt_cleanup | 30m | celery | → index_attempt_cleanup |
check_for_checkpoint_cleanup | 1h | celery | → checkpoint_cleanup |
check_for_hierarchy_fetching | 1h | celery | → connector_hierarchy_fetching |
monitor_background_processes | 5m | monitoring | — |
monitor_celery_queues | 10s | monitoring | 仅自托管 |
celery_beat_heartbeat | 1m | celery | 写 Redis 心跳键,仅自托管 |
check_for_doc_permissions_sync | 30s | celery | 仅 EE |
check_for_external_group_sync | 20s | celery | 仅 EE |
EE 两项的开关在 beat_schedule.py:201:ENTERPRISE_EDITION_ENABLED or _LICENSE_ENFORCEMENT_ENABLED。
4.2 三份清单,各有各的用途
beat_schedule.py 里并存三个列表,别混淆:
beat_task_templates(:38)——模板。云上会被make_cloud_generator_task(:302)转写成"每租户展开一次"的生成器任务;自托管则原样拷贝进tasks_to_schedule(:415-419)。beat_cloud_tasks(:329)——全系统级的云任务(监控 alembic、监控队列、检查可用租户),不按租户展开。tasks_to_schedule(:374)——自托管专属,只在not MULTI_TENANT时填充。
一个容易漏掉的细节:skip_gated / work_gated 这两个选项是云上专用提示,自托管路径会在拷贝时显式 pop 掉,免得它们泄漏成 apply_async 的非法参数(:411-419)。
4.3 expires 而不是堆积
每个任务都带 "expires": BEAT_EXPIRES_DEFAULT(15 分钟,:29)。注释讲得很直白:这些"检查"任务没必要 排队等,重要的是它们大致按点跑过。队列堵了 15 分钟以上的老 tick 直接作废,下一个 tick 会顶上——这避免了 Beat 在 worker 卡住时把队列灌爆。
5. 单例与栅栏(本章的核心机制)
分布式系统里最贵的一件事是**"同一件活别干两遍"**。Onyx 用了四种不同强度的手段,从强到弱依次是:
强 ┌──────────────────────────────────────────────┐
│ ① Postgres 行锁 nowait 索引去重(最新) │
│ ② Redis 分布式锁 + 续租 primary 单例 │
│ ③ Redis fence 键 + TTL 工作流占位(旧) │
│ ④ Redis beat lock 非阻塞 定时任务不重叠 │
弱 └─────── ───────────────────────────────────────┘
5.1 primary 的单例锁与 Bootstep 续租
要解决的小问题: 只能有一个 primary。它开机时会清空一堆 Redis 状态(primary.py:184-198),两个 primary 同时清就等于互相拆台。
做法: 开机时抢一把 Redis 锁,抢不到直接自杀。
# backend/onyx/background/celery/apps/primary.py:165-177(真实源码,节选)
lock: RedisLock = r.lock(
OnyxRedisLocks.PRIMARY_WORKER,
timeout=CELERY_PRIMARY_WORKER_LOCK_TIMEOUT,
thread_local=False,
)
acquired = lock.acquire(blocking_timeout=CELERY_PRIMARY_WORKER_LOCK_TIMEOUT / 2)
...
raise WorkerShutdown("Primary worker lock could not be acquired!")
锁超时 120 秒(backend/onyx/configs/constants.py:127 CELERY_PRIMARY_WORKER_LOCK_TIMEOUT),thread_local=False 是必须的——续租发生在另一个线程上,线程本地的 token 会让 reacquire() 认不出自己的锁。
续租为什么不能做成一个 Beat 任务? 答案就写在类的 docstring 里:
# backend/onyx/background/celery/apps/primary.py:269-275(真实源码,节选)
class HubPeriodicTask(bootsteps.StartStopStep):
"""Regularly reacquires the primary worker lock outside of the task queue.
...
This cannot be done inside a regular beat task because it must run on schedule and
a queue of existing work would starve the task from running.
"""
如果续租排在任务队列里,一旦队列积压,续租就排不上号 → 锁过期 → primary 以为自己不是 primary 了。所以它挂在 Celery Bootstep 上,直接用 worker 事件循环(hub)的定时器,绕开队列。续租间隔是 120 / 8 = 15 秒(:281),8 倍的安全余量。
时序上是这样:
t=0 抢到锁,TTL 120s
t=15 hub 定时器触发 → lock.owned()? → reacquire(),TTL 重置 120s
t=30 同上 …
↑ 连错 7 次才会真的过期
异常 lock.owned() == False → 重新 acquire(),日志 warning
5.2 Beat 的心跳与 watchdog:谁来看着看门人
Beat 没有单例锁,也没有 autorestart(supervisord.conf:126-132 没写 autorestart=true)。它一旦僵死,整个系统会安静地停止一切定时工作——不报错,只是不再有任何东西被调度。
Onyx 的解法是一条绕一圈的心跳链路:
Beat ──每 1 分钟调度──▶ primary 执行 celery_beat_heartbeat
│
▼
Redis SET onyx:celery:beat:heartbeat (TTL 600s)
│
每 60s 读一次 │
▼
supervisord_watchdog(独立 Python 进程)
│
连续 >5 次读不到 且 距上次成功 >900s
▼
supervisorctl restart celery_beat
三个锚点:任务在 backend/onyx/background/celery/tasks/shared/tasks.py:359-365(celery_beat_heartbeat,ex=600);键名常量 backend/onyx/configs/constants.py:723;watchdog 阈值 backend/onyx/utils/supervisord_watchdog.py:12-14(MAX_AGE_SECONDS = 900、CHECK_INTERVAL = 60、MAX_LOOKUP_FAILURES = 5)。supervisord 里那句 --key "onyx:celery:beat:heartbeat" 是手抄的字符串,配置文件里专门写了注释提醒它必须和常量一致(supervisord.conf:136)。
妙在哪:心跳不是 Beat 自己写的。Beat 只负责"调度",真正写键的是 primary。所以这条链路同时验证了"Beat 活着且它派的活真的能被执行"——Beat 假死(进程在、tick 不动)也能被抓到。
5.3 Redis fence:一次工作流的占位牌
要解决的小问题: 删除一个连接器会派生出成百上千个子任务。怎么知道"这场删除还在进行中"?
思路: 给整场工作流插一面旗(fence 键),旗在就代表流程未结束;旗上还写着 payload(有多少子任务)。
# backend/onyx/redis/redis_connector_delete.py:76-83(真实源码,节选)
def set_fence(self, payload: RedisConnectorDeletePayload | None) -> None:
if not payload:
self.redis.srem(OnyxRedisConstants.ACTIVE_FENCES, self.fence_key)
self.redis.delete(self.fence_key)
return
self.redis.set(self.fence_key, payload.model_dump_json(), ex=self.FENCE_TTL)
self.redis.sadd(OnyxRedisConstants.ACTIVE_FENCES, self.fence_key)
三层设计值得学:
| 键 | TTL | 作用 |
|---|---|---|
connectordeletion_fence_<id> | 7 天 | 旗本身,存 payload;TTL 纯属防内存泄漏 |
connectordeletion_active_<id> | 1 小时 | "还有人在动"的心跳信号,set_active() 刷新 |
active_fences(全局 SET) | — | 所有活跃 fence 的索引,避免 SCAN 全库 |
为什么要 active 这层?docstring 说得很清楚:"it's impossible to get the exact state of the system at a single point in time"(redis_connector_delete.py:88-92)——单看队列长度和任务状态会有竞态,需要一个带 TTL 的"最近还在动"信号来把时间上的缝补上。
primary 开机时会把这些全部推倒重来:r.delete(OnyxRedisConstants.ACTIVE_FENCES) 加上七个 reset_all()(primary.py:186-198)。