跳到主要内容

数据截至 (上游 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_serverFastAPI 应用,挂 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:442OnyxCeleryQueues-Q 绑定来自 backend/supervisord.conf,并发默认值来自 backend/onyx/configs/app_configs.py

worker消费的队列(-Q默认并发prefetch干什么
primarycelery41单例协调者:扫状态、派活、清 Redis
lightvespa_metadata_sync, connector_deletion, doc_permissions_upsert, checkpoint_cleanup, index_attempt_cleanup, opensearch_migration248短平快:元数据同步、删文档、清检查点
heavyconnector_pruning, connector_doc_permissions_sync, connector_external_group_sync, csv_generation, sandbox41长耗时、打外部 API:剪枝、权限同步、CSV 导出、沙箱
docfetchingconnector_doc_fetching11跑连接器拉原始文档,按批落盘
docprocessingdocprocessing61取批次跑索引管线(切块、嵌入、入库)
user_file_processinguser_file_processing, user_file_project_sync, user_file_delete21用户上传的文件
scheduled_tasksscheduled_tasks41Craft 定时 agent 执行器(长跑 LLM 调用)
monitoringmonitoring11队列深度、连接器成败、进程内存等指标

并发默认值锚点: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_migrationchat_ttl_deletionbackend/supervisord.conf:46-47
未提及 scheduled_tasks worker存在且有独立 app/config/supervisord programbackend/supervisord.conf:91-100

4. Beat:整个系统的节律

4.1 定时表

Beat 自己不干活,它只按点往队列里丢"去检查一下"的任务。自托管模式下的完整节律(backend/onyx/background/celery/tasks/beat_schedule.py:42-225):

任务周期落到哪个队列派生出什么
check_for_indexing15scelery(primary)connector_doc_fetching
check_for_vespa_sync_task20sceleryvespa_metadata_sync
check_for_pruning20sceleryconnector_pruning
check_for_connector_deletion20sceleryconnector_deletion
check_for_user_file_processing / _project_sync / _delete各 20sceleryuser_file_*
dispatch_due_scheduled_tasks30sceleryscheduled_tasks
check_for_index_attempt_cleanup30mceleryindex_attempt_cleanup
check_for_checkpoint_cleanup1hcelerycheckpoint_cleanup
check_for_hierarchy_fetching1hceleryconnector_hierarchy_fetching
monitor_background_processes5mmonitoring
monitor_celery_queues10smonitoring仅自托管
celery_beat_heartbeat1mcelery写 Redis 心跳键,仅自托管
check_for_doc_permissions_sync30scelery仅 EE
check_for_external_group_sync20scelery仅 EE

EE 两项的开关在 beat_schedule.py:201ENTERPRISE_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-365celery_beat_heartbeatex=600);键名常量 backend/onyx/configs/constants.py:723;watchdog 阈值 backend/onyx/utils/supervisord_watchdog.py:12-14MAX_AGE_SECONDS = 900CHECK_INTERVAL = 60MAX_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)。

5.4 索引的去 Redis 化:IndexingCoordination

Onyx 正在从 Redis fence 迁往 Postgres 行锁。 索引这条路径已经迁完,类的 docstring 直说:"Database-based coordination for indexing tasks, replacing Redis fencing"backend/onyx/db/indexing_coordination.py:33-34)。

先看它想达到的效果(示意):

# 示意,非源码 —— 「同一个 cc_pair 不能有两个全量索引」的核心想法
def try_start(cc_pair_id, search_settings_id):
# 1) 排他锁住"该 cc_pair 上还没结束的全量任务"这几行,锁不到就立刻放弃
existing = select_active_attempts(cc_pair_id, search_settings_id).for_update(nowait=True)
if existing:
return None # 已有人在跑 → 本次不派活
return create_index_attempt(...) # 插入这一行,就等于"立了旗"

真实实现的关键三行:

# backend/onyx/db/indexing_coordination.py:59-64(真实源码,节选)
IndexAttempt.status.in_(
[IndexingStatus.NOT_STARTED, IndexingStatus.IN_PROGRESS]
),
IndexAttempt.targeted_reindex_job_id.is_(None),
...
.with_for_update(nowait=True)

三个决策点,逐个说:

  • nowait=True——抢不到锁就立即SQLAlchemyError,被 :95-103 的 except 吞掉并返回 None。这正是我们要的语义:并发到来时"退让"而不是"排队"。
  • targeted_reindex_job_id.is_(None)——定向重索引任务被故意排除在这个 fence 之外,允许它和全量爬取重叠;冲突交给逐文档的行锁处理(注释在 :50-53)。
  • "创建 IndexAttempt 行 = 立旗":76 的注释 # Create new index attempt (this is setting the "fence"))——旗和事实合一,没有第二份状态可以不一致。

相应地,取消也不再走 Redis 信号,而是把 cancellation_requested 写进同一行(:117-131 request_cancellation:106-115 check_cancellation_requested),由长跑任务轮询。

5.5 两段式派发:docfetching → docprocessing

索引的实际执行被拆成两跳,中间用文件存储而不是消息体传批次:

Beat(15s)


check_for_indexing docprocessing/tasks.py:847
│ 先抢 CHECK_INDEXING_BEAT_LOCK(非阻塞,抢不到直接 return)

try_creating_docfetching_task docfetching/task_creation_utils.py:20
│ ① 抢 "try_creating_indexing_task" 函数级锁(30s)
│ ② IndexingCoordination.try_create_index_attempt ← 真正的去重点
│ ③ send_task → connector_doc_fetching 队列

docfetching worker
│ 另起进程跑连接器;每攒够一批文档:
│ 批次写入 file store,只把 batch_num 放进消息

send_task → docprocessing 队列 run_docfetching.py:788-793

docprocessing worker ── 切块 / 嵌入 / 入索引(见第 05 章)

三道闸门层层收窄,值得留意它们各自防的是什么

闸门位置防什么
CHECK_INDEXING_BEAT_LOCK(非阻塞)docprocessing/tasks.py:895-902两次 beat tick 的 check_for_indexing 重叠
try_creating_indexing_task 函数锁(30s)task_creation_utils.py:42-50定时触发与手动 API 触发撞车
try_create_index_attempt 行锁indexing_coordination.py:66同一 cc_pair 起两个全量索引

还有个小而实用的优先级技巧:首次索引的连接器用 HIGH,已经成功过的用 MEDIUMtask_creation_utils.py:77-83)。新接入的连接器不会被存量连接器的重索引饿死。


6. 多租户:一份代码,N 个 schema

6.1 一个开关

# backend/shared_configs/configs.py:178, 183, 197(真实源码,节选)
MULTI_TENANT = os.environ.get("MULTI_TENANT", "").lower() == "true"
POSTGRES_DEFAULT_SCHEMA = os.environ.get("POSTGRES_DEFAULT_SCHEMA") or "public"
TENANT_ID_PREFIX = "tenant_"

关掉时,"租户 id"恒等于 public,所有多租户代码路径退化成直连;打开时,一个租户 = 一个 Postgres schema,名字形如 tenant_<uuid>

6.2 隔离靠的是 contextvar + schema 翻译

隔离不是靠每个查询里手写 WHERE tenant_id = ?,而是靠两段配合:

第一段:一个进程级的 contextvar 记住"我现在是谁"。

# backend/shared_configs/contextvars.py:7-9(真实源码,节选)
CURRENT_TENANT_ID_CONTEXTVAR: contextvars.ContextVar[str | None] = (
contextvars.ContextVar(
"current_tenant_id", default=None if MULTI_TENANT else POSTGRES_DEFAULT_SCHEMA
)

注意这个 default 的巧思:自托管时它天生就是 public,所以没有任何一处需要"记得设置租户"。

第二段:开会话时把 schema 名交给 SQLAlchemy,由它改写所有 SQL。

# backend/onyx/db/engine/sql_engine.py:464-466(真实源码,节选)
schema_translate_map = {None: tenant_id}
with engine.connect().execution_options(
schema_translate_map=schema_translate_map
) as connection:

所有 ORM 模型都不写 schema(即 None),运行时统一被翻译成当前租户的 schema。业务代码一行都不用改就获得了隔离。

自托管 + 默认 schema 时还有一条快路径,跳过翻译直接建 Session(sql_engine.py:531-538)。

6.3 TenantAwareTask:worker 侧的接力棒

问题是:contextvar 无法跨进程传递。Celery 消息到了 worker 那头,谁来设它?答案是把这件事塞进所有任务的基类

# backend/onyx/background/celery/apps/app_base.py:84-97(真实源码,节选)
def __call__(self, *args: Any, **kwargs: Any) -> Any:
tenant_id = kwargs.get("tenant_id", None) or POSTGRES_DEFAULT_SCHEMA
CURRENT_TENANT_ID_CONTEXTVAR.set(tenant_id)
try:
return super().__call__(*args, **kwargs)
finally:
CURRENT_TENANT_ID_CONTEXTVAR.set(None)

finally 里的复位是安全关键:线程池 worker 会在同一个线程上连续执行不同租户的任务,不清就会串租户。每个 app 都要挂:celery_app.Task = app_base.TenantAwareTask(如 primary.py:55),Beat 侧则是 celery_app.conf.task_default_base = app_base.TenantAwareTaskbeat.py:275)。

于是租户 id 的传递链是:Beat 把 tenant_id 写进 kwargs → 消息体 → TenantAwareTask.__call__ 读出来设 contextvar → 数据库会话据此翻译 schema

6.4 Redis 侧:key 前缀隔离

Redis 没有 schema,所以 Onyx 自己造了一层。TenantRedisClient组合而非继承重写每个碰 key 的方法,逐个显式加前缀:

# backend/onyx/redis/tenant_redis_client.py:1-8(真实源码,docstring 节选)
"""Replaces the old ``__getattribute__``-based ``TenantRedis`` (a ``redis.Redis``
subclass that wrapped a hand-maintained allowlist of methods) with an explicit,
hand-written client built by composition. ... Calling a Redis method that is not
exposed here is a typing error, not a silent cross-tenant write."""

这是本章最值得抄走的一条设计取舍:从"黑名单式的动态拦截"改成"白名单式的显式包装",把跨租户写入从运行时的静默 bug 变成编译期的类型错误。前缀函数还刻意做成幂等的(:56-57:已带前缀就原样返回),避免双重加前缀。

6.5 两套 alembic

backend/alembic.ini 里有两个 section:

sectionscript_location管什么
[alembic]alembic租户 schema 内的表;自托管时就是 public
[schema_private]alembic_tenants共享 public 表,仅多租户时使用(租户注册表等)

依据:backend/alembic.ini:112-118alembic_tenants/README.md 一句话说明:"public table migrations when operating with multi tenancy"。

租户版迁移的批量执行能力做得相当细:alembic -x upgrade_all_tenants=true-x schemas=...-x tenant_range_start/end=...(按字母序切片,方便分批灰度)、-x continue=true(出错继续)——全部记在 backend/alembic/README.md。配套的 get_schemas_needing_migrationbackend/onyx/db/engine/tenant_utils.py:41)用服务端 PL/pgSQL 循环逐个 schema 取版本号,避免拼出一条巨大的 UNION ALL

6.6 DynamicTenantScheduler:会自己长大的 Beat

标准 Celery Beat 的 schedule 是启动时静态确定的。云上租户随时新增,这不行。

DynamicTenantScheduler 继承 PersistentScheduler,在每次 tick() 里检查是否到了 60 秒的重载间隔(beat.py:31 RELOAD_INTERVAL = 60:61-80 tick),到了就:

_try_updating_schedule() beat.py:155

├─ get_all_tenant_ids() tenant_utils.py:120
│ 多租户:查 information_schema.schemata 列出所有 schema
│ 自托管:直接返回 [POSTGRES_DEFAULT_SCHEMA]

├─ OnyxRuntime.get_beat_multiplier() 可在 Redis 里在线调速

├─ _generate_schedule(tenants, multiplier) beat.py:82
│ 云: cloud 任务挂一次 + 每租户展开模板,周期 × multiplier
│ 自托管: 只展开一个 "public" 租户

└─ 只比较「任务名集合」是否变化 → 变了才 clear+update+sync

三个务实的细节:

  • 只比名字,不比周期_compare_schedules:240-246)。周期变化由 beat_multiplier != last_beat_multiplier 这条独立分支捕获(:177-179)。比全量 schedule 的深比较便宜得多。
  • beat_multiplier 默认 8.0beat_schedule.py:34),注释直说是 "hack to slow down task dispatch in the cloud until we have a better implementation (backpressure, etc)"。云上 check_for_indexing 的真实周期是 15s × 8 = 120s。诚实的临时方案。
  • 构造函数里不碰数据库beat.py:52-53 的注释),首次 schedule 生成推迟到 beat_init 信号、引擎初始化之后(:249-264)。

7. API 层的装配

7.1 一条流水线

get_application()backend/onyx/main.py:492)不是一段配置,而是一条有顺序要求的装配线:

FastAPI(...) ← openapi/docs 路由受 ENABLE_PUBLIC_DOCS 控制,默认不注册

├─ 异常处理器(400/401/403/404/500 + 领域异常) :498-506

├─ 60+ 业务 router :508-579
│ 全部经 include_router_with_global_prefix_prepended

├─ 认证 router(按 AUTH_TYPE 分支挂载) :581-728

├─ 中间件(注册顺序 = 执行顺序的逆序!) :742-773
│ CORS → Captcha ×2 → ClientIP → request-id → prometheus

├─ check_router_auth(application) :776
│ ★ 必须在所有路由注册完之后:逐个检查每个端点
│ 要么有鉴权依赖,要么被显式标记为 public

└─ use_route_function_names_as_operation_ids :778

统一前缀这件小事被抽成了一个函数,好处是"全局 API 前缀"只有一处实现:

# backend/onyx/main.py:242-248(真实源码,节选)
"""Adds the global prefix to all routes in the router."""
processed_global_prefix = f"/{APP_API_PREFIX.strip('/')}" if APP_API_PREFIX else ""
passed_in_prefix = cast(str | None, kwargs.get("prefix"))
if passed_in_prefix:
final_prefix = f"{processed_global_prefix}/{passed_in_prefix.strip('/')}"

check_router_auth:776)是我认为最值得抄的一行:把"忘记加鉴权"从代码审查问题变成启动期崩溃。新加一个端点忘了写依赖,服务直接起不来。

7.2 认证栈:一个 AUTH_TYPE 决定挂哪些路由

Onyx 用 fastapi-users 作底座,按 AUTH_TYPE 分支组装:

AUTH_TYPE挂载的 router锚点
BASIC / CLOUD登录、注册、改密、验证、用户管理(全套 fastapi-users)main.py:581-608
BASIC / CLOUD / GOOGLE_OAUTH移动端 bearer 网关 /auth/mobile:613-618
GOOGLE_OAUTH(或 BASIC + 配了凭证)Google OAuth2 web 版 + 一份独立的移动版:622-670
OIDC通用 OpenID router(强制补 offline_access scope 以拿 refresh token):672-708
SAMLsaml_routerEE 提供:710-714
DISABLEDrefresh token router:716-728

移动端为什么要另起一份 OAuth router?注释解释得很具体:回调地址必须落在 /api 下让 api_server 直接接住,走 web 的回调包装器会丢掉那个无 cookie 的深链 302(:648-651)。

7.3 两种限流,别搞混

维度HTTP 次数限流Token 预算限流
挡什么请求频次LLM token 消耗量
实现fastapi-limiter + Redis查 Postgres 里的用量记录
挂载方式router 级 Depends端点级 Depends
位置backend/onyx/server/middleware/rate_limiting.pybackend/onyx/server/query_and_chat/token_limit.py
IP + User-Agent,或用户 id全局 / 用户 / 用户组

次数限流有两组:认证端点用 IP+UA 作键(rate_limiting.py:36-42),反馈端点用用户 id 作键(:46-53)。后者的 docstring 把威胁模型写清楚了:以用户 id 为键,攻击者就没法靠"换 cookie / 换 Authorization 格式"给同一个账号刷出多个限流桶;匿名用户则退回 IP+UA,免得一个匿名滥用者耗光所有匿名者的额度(:72-87)。

Token 限流是本章第 8 节那个可插拔机制的教科书示例

# backend/onyx/server/query_and_chat/token_limit.py:33-41(真实源码,节选)
if not any_rate_limit_exists():
return
versioned_rate_limit_strategy = fetch_versioned_implementation(
"onyx.server.query_and_chat.token_limit", _check_token_rate_limits.__name__
)
return versioned_rate_limit_strategy(user)

MIT 版的 _check_token_rate_limits 只查全局限额(:44-45);EE 版把全局、用户、用户组三条并行查(ee/onyx/server/query_and_chat/token_limit.py:35-51,用 run_functions_tuples_in_parallel)。挂载点只有一处:chat_backend.py:801_rate_limit_check: None = Depends(check_token_rate_limits)


8. 开源分层:MIT 主干 + 镜像式 EE

8.1 目录是镜像的

backend/onyx/(MIT)与 backend/ee/onyx/(企业许可)同名同构

backend/
├── onyx/ ← MIT,完整可用
│ ├── auth/ background/ db/ server/ main.py …
│ └── LICENSE 覆盖:仓库根 MIT
└── ee/
├── LICENSE ← Onyx Enterprise License
└── onyx/
├── auth/ background/ db/ server/ main.py …
└── 只放"覆盖或新增"的那部分,不重复 MIT 代码

许可边界写在仓库根 LICENSE:5-8所有位于 "ee" 目录下的内容按 Onyx Enterprise License 授权,共三处(backend/eeweb/src/app/eeweb/src/ee)。

ee/onyx/main.py:81-93 就是这套镜像的缩影——它不重写应用,而是拿 MIT 的 get_application_base(...) 再往上叠中间件(tier gate、租户追踪或许可强制)。

8.2 三个函数,三种失败语义

接缝只有三个函数(backend/onyx/utils/variable_functionality.py),区别在找不到 EE 实现时怎么办

函数行号EE 缺失时典型用途
fetch_versioned_implementation:70回落到 MIT 同名符号,实在没有就抛有 MIT 默认行为、EE 想加强的(限流策略、get_application
fetch_versioned_implementation_with_fallback:125返回调用方给的 fallback允许静默降级的场景
fetch_ee_implementation_or_noop:161返回一个吞掉参数的 no-op(同步/异步自适应)纯 EE 能力,社区版就是不该发生

fetch_ee_implementation_or_noop 是本章标题里那个"可插拔企业实现"模式的主角。它的关键在于返回的是一个可调用对象,而不是一个值

# 示意,非源码 —— 调用方无需任何 if/else
censored = fetch_ee_implementation_or_noop(
"onyx.external_permissions.post_query_censoring", # EE 模块路径
"_post_query_chunk_censoring", # EE 符号名
retrieved_chunks, # 社区版的"什么都不做"返回值
)(retrieved_chunks, user) # ← 无论哪种情况,调用姿势完全一致

真实调用点:backend/onyx/context/search/pipeline.py:339-343(查询后按外部权限二次删除结果)、backend/onyx/access/access.py:148-151(索引时是否顺带抓权限)、backend/onyx/main.py:373-375(许可席位指标)。全仓 91 处 fetch_ee_implementation_or_noop + 48 处 fetch_versioned_implementation

这个模式的代价:符号是字符串,静态检查抓不到。EE 侧改名而 MIT 侧的字符串没跟着改,社区版会静默变 no-op,企业版则在 :191-197 那里抛错——单元测试之外没有别的护栏。

8.3 后台任务侧的挂接方式

EE 的 Celery 任务不改 supervisord 命令,而是在 ee/onyx/background/celery/apps/primary.py导入 MIT 的 app 再追加 autodiscover

# backend/ee/onyx/background/celery/apps/primary.py:1-14(真实源码,节选)
from onyx.background.celery.apps.primary import celery_app
celery_app.autodiscover_tasks(
app_base.filter_task_modules([
"ee.onyx.background.celery.tasks.hooks",
"ee.onyx.background.celery.tasks.doc_permission_syncing",
"ee.onyx.background.celery.tasks.external_group_syncing",
"ee.onyx.background.celery.tasks.cloud",
...
])
)

Beat 侧同理:ee/onyx/background/celery/tasks/beat_schedule.py 导入 MIT 的三份清单,追加 EE 条目(用量报告、TTL 管理、查询历史导出清理等),再重新导出同名的 get_tasks_to_schedule / get_cloud_tasks_to_schedule——正好被 beat.py:117-119fetch_versioned_implementation 接住。MIT 侧完全不知道 EE 的存在

8.4 EE-only 能力清单(以本 commit 的目录与调用点为准)

能力位置
SAML 登录main.py:710-714saml_router 来自 EE
外部权限同步 / 查询后审查ee/onyx/external_permissions/ee/.../doc_permission_syncing
外部用户组同步ee/.../external_group_syncing
用户组 RBACee/onyx/server/user_group/
分析、查询历史导出、用量报告`ee/onyx/server/analytics
许可强制、席位计量、tier gateee/onyx/server/license/ee/onyx/server/middleware/{license_enforcement,tier_gate}.py
多租户控制面(租户注册、计费、SCIM)ee/onyx/server/{tenants,billing,scim}/
按用户 / 用户组的 token 限流ee/onyx/server/query_and_chat/token_limit.py

9. 部署形态:Lite 与 Standard

9.1 两档

README(README.md:72-86)给出两种:

Onyx LiteStandard Onyx
定位轻量 Chat UI + Agents完整功能
内存< 1GB显著更高
向量索引无(DISABLE_VECTOR_DB=trueOpenSearch / Vespa
后台 worker不启动全套 Celery
模型推理服务不启动inference + indexing 两份
缓存 / 会话PostgresRedis
文件存储PostgresS3 / MinIO

9.2 Lite 是一层 compose overlay,不是另一套代码

docker compose -f docker-compose.yml -f docker-compose.onyx-lite.yml up -d

overlay 干的事只有两类(deployment/docker_compose/docker-compose.onyx-lite.yml:17-31):把 backgroundcache、两个 model server、opensearchminio 全部塞进 profile(默认不启动),并把 api_server 的 depends_on 改成 required: false;再设四个环境变量把后端切到 Postgres(:51-55)。想找回任何一件,加对应的 --profile 即可。

Lite 下谁干后台的活? 一个从 API 服务器里长出来的守护线程:

# backend/onyx/background/periodic_poller.py:1-12(真实源码,docstring 节选)
"""Periodic poller for NO_VECTOR_DB deployments.

Replaces Celery Beat and background workers with a lightweight daemon thread
that runs from the API server process. ..."""

它做两件事:每 30 秒把卡在 PROCESSING/DELETING 的用户文件重新推一遍;按配置的间隔跑 LLM 模型更新与定时评测,用 Postgres advisory lock 在多个 API 实例之间去重:10-12:26 PERIODIC_TASK_LOCK_BASE)。这是 Redis 缺席时"锁"的替代方案。

Standard 的 background 容器则是一个 supervisord_entrypoint.shdocker-compose.yml:125-130),一口气拉起前面表里那 8 个 worker + Beat + watchdog + Slack/Discord bot。

9.3 Kubernetes:worker 各自成 Deployment

Helm chart 把每种 worker 拆成独立 Deployment,并各配 HPA / KEDA ScaledObject / metrics Service:deployment/helm/charts/onyx/templates/celery-worker-{primary,light,heavy,docfetching,docprocessing,monitoring,scheduled-tasks,user-file-processing}*.yaml-Q 参数与 supervisord 里逐字一致(如 celery-worker-light.yaml:74-76,两边现在都是 7 个队列)。

探针文件由 Celery Bootstep 写:LivenessProbe 每 15 秒 touch 一次 /tmp/onyx_k8s_<worker>_liveness.txt,停止时删除(app_base.py:699-722celery_utils.py:269-283 make_probe_path)。Beat 的探针在 tick() 里顺带 touch(beat.py:69)。

9.4 model_server 的边界

backend/model_server/唯一一个不 import 业务代码的后端目录。它只是个 FastAPI 应用,挂两个 router:管理端点与编码端点(model_server/main.py:132-133)。同一个镜像按 INDEXING_ONLY 环境变量跑成两份实例,区别仅在日志前缀(IDX vs INF:129-132)和资源配额。

这条边界的好处很直接:GPU 只需要给 model_server,业务进程可以纯 CPU 横向扩;模型热加载不影响 API 可用性。


10. 边界与已知薄弱点

诚实清单。这些都是按本 commit 的代码读出来的。

10.1 Redis 挂了会怎样

Redis 是协调中枢,不是可选缓存。硬依赖它存活的路径:

路径后果
Celery broker所有后台任务停摆
primary 单例锁续租锁过期 → 后续行为不确定
Beat 心跳 → watchdog读不到键,watchdog 会误判 Beat 死了并重启它
各类 fence / active 键删除、剪枝、权限同步的进度信息丢失
停止生成的 fence停止按钮失效(除非切 Postgres 缓存后端)

worker 启动时会等 Redis,60 秒等不到就 WorkerShutdown 自杀(app_base.py:327-366WAIT_LIMIT = 60),交给 supervisord 重启——快速失败而非降级运行

唯一的逃生舱是 Lite 模式:CACHE_BACKEND=postgres + AUTH_BACKEND=postgres 能让 Redis 完全缺席(app_configs.py:85-87:169),代价是没有后台 worker。

10.2 primary 的"单例"是半成品

代码注释自己承认了:

# backend/onyx/background/celery/apps/primary.py:153-161(真实源码,注释节选)
# For the moment, we're assuming that we are the only primary worker
# that should be running.
# TODO: maybe check for or clean up another zombie primary worker if we detect it
r.delete(OnyxRedisLocks.PRIMARY_WORKER)
...
# it is planned to use this lock to enforce singleton behavior on the primary
# worker, since the primary worker does redis cleanup on startup, but this isn't
# implemented yet.

注意 r.delete(PRIMARY_WORKER) 这一句:它先把锁删掉再抢。所以此刻的锁并不能阻止第二个 primary 起来——真起来了,它会把前一个的锁抹掉、把 Redis 状态清空。防重复部署目前靠的是编排层(supervisord 只配一份、Helm 里 replicas=1),不是这把锁。

10.3 多租户下开机自检被跳过

# backend/onyx/background/celery/apps/primary.py:131-133(真实源码,节选)
# Less startup checks in multi-tenant case
if MULTI_TENANT:
return

也就是说:云上不做 Redis 状态清理、不抢单例锁、不标记孤儿索引任务。这些在自托管下由 primary 开机完成的收尾工作,在多租户下要么由云控制面另行处理,要么就是缺失的——从这份克隆的代码里看不出替代实现在哪。

10.4 一个没有消费者的队列

connector_hierarchy_fetchingconstants.py:434 有定义,hierarchyfetching/tasks.py:144 会往它派任务,monitoring/tasks.py:174 会统计它的长度——但它不出现在 backend/supervisord.conf 的任何 -Q 里,也不出现在任何 Helm 模板的 -Q。heavy worker 的 app 确实 autodiscover 了 hierarchyfetching 任务模块(apps/heavy.py:136),但没订阅那个队列。按默认部署配置,check_for_hierarchy_fetching 派出去的任务会一直躺在队列里 (inferred)。

10.5 停止生成:只能丢弃后续输出,无法真正中断线程

这是本章最该讲清楚的一个"物理限制"。停止按钮的实现是一面 fence:

# backend/onyx/chat/stop_signal_checker.py:38-48(真实源码,节选)
def is_connected(chat_session_id: UUID, cache: CacheBackend) -> bool:
"""Check if the chat session should continue (not stopped)."""
return not cache.exists(_get_fence_key(chat_session_id))

流式主循环每 50 毫秒(process_message.py:1120 _CANCEL_POLL_INTERVAL_S = 0.05)在队列 get 超时的间隙检查一次这面旗。发现旗立起来之后:

# backend/onyx/chat/process_message.py:1385-1393(真实源码,节选)
_publish(Packet(placement=Placement(turn_index=0),
obj=OverallStop(type="stop", stop_reason="user_cancelled")))
# Workers can't be interrupted; discard their remaining
# output via the Emitter no-op.
drain_done.set()
return

注释就是答案:Workers can't be interrupted。Python 没有安全的线程强杀原语,正在跑的模型调用线程会一直跑到自然结束。所以"停止"的真实语义是三件事——

  1. 立刻给前端发一个 OverallStop,UI 上生成停住;
  2. 把每个模型的当前输出快照持久化(_PersistContext.STOP_BUTTON:1379-1382),已完成的存完整、在飞的存部分并标记 stopped;
  3. 把 Emitter 变成 no-op,后续 worker 吐出来的东西全部丢弃

代价诚实地摆着:token 已经付了;如果调用的是外部 LLM API,请求还在跑;worker 线程会在后台占着资源直到自己结束(:1834 附近的注释也提到"Worker threads may still be running (e.g. user-cancellation path)")。

fence 的 TTL 是 10 分钟(stop_signal_checker.py:7),下一轮开始前由 reset_cancel_status 清掉(process_message.py:1120)。

10.6 其它

  • 文档漂移backend/onyx/background/README.md 的 worker 表已滞后于代码(见 §3.4)。
  • EE 接缝无静态检查:模块路径与符号名都是字符串(§8.2)。
  • docprocessing 有意不开 task_track_started,配置注释说明原因:"celery occasionally runs tasks more than once",开了反而会让任务状态被重复运行改乱(configs/docprocessing.py:22-28)。这句话本身也是一条重要信息——Onyx 假设 Celery 任务可能被执行多次,去重必须靠 §5 的那几道闸门,不能靠"任务只跑一次"。

11. 代码地图(导航索引)

主题文件路径符号名
worker→queue 权威映射backend/supervisord.conf[program:celery_worker_*]-Q
worker 分工文档(略有滞后)backend/onyx/background/celery/README.mdWorker → Queue Mapping
队列名常量backend/onyx/configs/constants.pyOnyxCeleryQueues
锁名常量backend/onyx/configs/constants.pyOnyxRedisLocks
并发默认值backend/onyx/configs/app_configs.pyCELERY_WORKER_*_CONCURRENCY
各 worker 的 Celery 配置backend/onyx/background/celery/configs/*.pyworker_concurrency / worker_pool
EE/MIT 双版本入口壳backend/onyx/background/celery/versioned_apps/*.pyfetch_versioned_implementation
primary 单例锁与开机清理backend/onyx/background/celery/apps/primary.pyon_worker_init
锁续租 Bootstepbackend/onyx/background/celery/apps/primary.pyHubPeriodicTask
定时节律定义backend/onyx/background/celery/tasks/beat_schedule.pybeat_task_templates, get_tasks_to_schedule
云端每租户任务生成backend/onyx/background/celery/tasks/beat_schedule.pymake_cloud_generator_task
动态租户调度器backend/onyx/background/celery/apps/beat.pyDynamicTenantScheduler
Beat 心跳任务backend/onyx/background/celery/tasks/shared/tasks.pycelery_beat_heartbeat
Beat 看门狗backend/onyx/utils/supervisord_watchdog.pymain, MAX_AGE_SECONDS
Redis fence 三层键backend/onyx/redis/redis_connector_delete.pyset_fence, set_active, fenced
索引去重(DB 行锁)backend/onyx/db/indexing_coordination.pyIndexingCoordination.try_create_index_attempt
索引取消backend/onyx/db/indexing_coordination.pyrequest_cancellation
索引调度入口backend/onyx/background/celery/tasks/docprocessing/tasks.pycheck_for_indexing
索引任务创建backend/onyx/background/celery/tasks/docfetching/task_creation_utils.pytry_creating_docfetching_task
抓取→处理的批次派发backend/onyx/background/indexing/run_docfetching.pyOnyxCeleryTask.DOCPROCESSING_TASK
租户上下文变量backend/shared_configs/contextvars.pyCURRENT_TENANT_ID_CONTEXTVAR
租户任务基类backend/onyx/background/celery/apps/app_base.pyTenantAwareTask
schema 级隔离backend/onyx/db/engine/sql_engine.pyget_session_with_tenant, schema_translate_map
租户枚举与校验backend/onyx/db/engine/tenant_utils.pyget_all_tenant_ids, validate_tenant_id
Redis 租户前缀隔离backend/onyx/redis/tenant_redis_client.pyTenantRedisClient, _prefix_key
两套迁移的分工backend/alembic.ini[alembic], [schema_private]
FastAPI 应用装配backend/onyx/main.pyget_application, include_router_with_global_prefix_prepended
启动期鉴权自检backend/onyx/main.pycheck_router_auth
HTTP 次数限流backend/onyx/server/middleware/rate_limiting.pyget_auth_rate_limiters, get_feedback_rate_limiters
Token 预算限流backend/onyx/server/query_and_chat/token_limit.pycheck_token_rate_limits
EE 可插拔接缝backend/onyx/utils/variable_functionality.pyfetch_ee_implementation_or_noop
EE 应用装配backend/ee/onyx/main.pyget_application
EE 后台任务挂接backend/ee/onyx/background/celery/apps/primary.pyautodiscover_tasks
Lite 模式 overlaydeployment/docker_compose/docker-compose.onyx-lite.ymlprofiles, DISABLE_VECTOR_DB
Lite 模式的后台替身backend/onyx/background/periodic_poller.py_PeriodicTaskDef
模型服务进程backend/model_server/main.pyget_model_app, INDEXING_ONLY
停止生成的 fencebackend/onyx/chat/stop_signal_checker.pyset_fence, is_connected
停止生成的丢弃逻辑backend/onyx/chat/process_message.py_CANCEL_POLL_INTERVAL_S, OverallStop

相关章节: 一次对话的生命周期(本章的 api_server 里到底跑了什么)· 索引栈docfetching / docprocessing 两个 worker 内部的管线)· 检索栈model_server 被谁调用)