跳到主要内容

数据截至 (上游 commit 5053c08115bd)

索引栈:60+ 连接器如何变成可检索的 chunk

30 秒导读: 04 章讲的是"一句提问怎么变成带引用的答案"——那是读取侧。这一章讲写入侧:Slack 的一条消息、Confluence 的一个页面、GDrive 的一份 PDF,是怎么被抓下来、切成 chunk、算成向量、连同 ACL 一起写进向量库的;以及它们过期、被删、权限变了之后怎么被清理。


1. 这一章讲什么(零基础也能懂)

一句话定义: 索引栈 = 把"外部系统里的东西"变成"向量库里可被检索的 chunk"的那条流水线,外加维护这些 chunk 一直和源头保持一致的那些后台任务。

写入侧和读取侧是同一个索引的两端,职责刚好互补:

读取侧(04 章)写入侧(本章)
触发者用户提问定时任务 / 手动触发
输入一句自然语言外部系统的 API 响应
输出带引用的答案向量库里的 chunk 行
延迟预算秒级小时级,可以很慢
失败代价这次答不好数据缺失,用户永远搜不到

为什么这件事难? 三个各自独立的难点:

  • 源头不可控。 53 个数据源、53 套分页语义和限流规则,有的能给"最近变更",有的只能全量重扫。
  • 一次跑很久。 一个大 Confluence 站点跑几小时很正常,中途进程被杀、限流、单个页面 403 都是常态——不能"跑失败就从头再来"。
  • 数据要一直对得上。 源头删了文档、改了权限、你换了嵌入模型——索引都得跟着变,否则用户会搜到本不该看见的东西。

一句话直觉: 把它当成一条带存档点的 ETL 流水线。连接器是"外部世界的适配头",checkpoint 是存档点,ConnectorFailure 是"这一格坏了但游戏继续"的记录方式。

规模的真实数字: DocumentSource 枚举有 57 个值,索引连接器注册表 CONNECTOR_CLASS_MAP 里有 53 条映射、指向 50 个不同的连接器类(backend/onyx/connectors/registry.py:14);四种对象存储(S3/R2/GCS/OCI)共用一个 BlobStorageConnector。标题里的"60+"是把联邦检索连接器(backend/onyx/federated_connectors/)之类也算进去的口径。


2. 顶层全景(它大概怎么转)

2.1 一条文档的五站

怎么读这张图:从左到右是一条文档的旅程,每一站都可能把它丢掉或标记为失败,但不会中断整批。

外部系统 ①抓取 ②准备 ③加工 ④落库
┌────────┐ ┌───────────┐ ┌───────────┐ ┌───────────┐ ┌───────────┐
│ Slack │ │ 连接器 │ │ 去重 + 落 │ │ 切块 │ │ 拼 ACL │
│ Confl. │ ──▶ │ 吐 Document│──▶│ Postgres │──▶│ 补上下文 │──▶│ 写向量库 │
│ GDrive │ │ 吐 Failure │ │ 元数据 │ │ 算 embed │ │ 回写计数 │
└────────┘ └───────────┘ └───────────┘ └───────────┘ └───────────┘
│ │
└──── checkpoint 存档 ────┐ ┌───────┘
▼ ▼
Postgres(真源) 向量库(可检索)

2.2 部件一句话职责

部件干什么在哪个文件
连接器契约定义"一个数据源必须能提供什么"backend/onyx/connectors/interfaces.py
注册表 + 工厂DocumentSource 懒加载并实例化连接器backend/onyx/connectors/registry.pyfactory.py
ConnectorRunner把四种连接器风格统一成一个 generator,负责成批和存档backend/onyx/connectors/connector_runner.py
索引管线去重 → 图片摘要 → 切块 → 补上下文 → 嵌入 → 写库backend/onyx/indexing/indexing_pipeline.py
Chunker按 token 预算切块,把标题/元数据拼进可检索文本backend/onyx/indexing/chunker.py
DefaultIndexingEmbedder调独立 model server 把 chunk 文本变成向量backend/onyx/indexing/embedder.py
适配器把"写 Postgres/加 ACL"这些副作用从管线里剥出去backend/onyx/indexing/adapters/document_indexing_adapter.py

2.3 两个进程,一个 filestore 中转

关键结构决策:抓取和加工被拆成两个 Celery worker,中间靠文件存储 + Redis 队列解耦。

docfetching worker docprocessing worker
┌────────────────────┐ ┌────────────────────┐
│ ConnectorRunner.run│ 存一批 → filestore │ 取一批 ← filestore │
│ ↓ 每 16 篇一批 │ ══════════════════▶ │ ↓ │
│ store_batch(n) │ 发一个 celery 任务 │ run_indexing_ │
│ send_task(n) │ │ pipeline(...) │
│ save_checkpoint() │ │ │
└────────────────────┘ └────────────────────┘
慢、受源头限流 CPU/GPU 密集

抓取循环见 backend/onyx/background/indexing/run_docfetching.py:441connector_document_extraction:它在 :762 把清洗后的批次写进 batch_storage:788app.send_task(OnyxCeleryTask.DOCPROCESSING_TASK, ...) 派发,然后 :815 保存 checkpoint。加工侧在 backend/onyx/background/celery/tasks/docprocessing/tasks.py:1853run_indexing_pipeline

这样拆的收益:连接器被源头限流时,不会占着嵌入 worker;嵌入慢时,也不会拖住抓取。worker 分工的全貌见 06 章


3. 连接器契约:一份接口,53 个数据源

它要解决的小问题: 每个数据源的 API 都不一样,但下游管线只想拿到"一批 Document"。

思路: 不做一个万能大接口,而是拆成一组小的抽象基类,连接器按自己源头的能力挑着实现。instantiate_connector 出来的对象是什么类型,ConnectorRunner 就用 isinstance 走哪条路径。

3.1 四类主流:按"能不能增量、能不能续跑"分

契约核心方法语义谁在用它
LoadConnectorload_from_state()全量:把源头当前完整状态吐一遍首次索引 / 无法增量的源
PollConnectorpoll_source(start, end)增量:只吐这个时间窗内变过的有"最近修改"API 的源
SlimConnectorretrieve_all_slim_docs()只吐 ID,不吐正文剪枝任务(见 §7.1)
CheckpointedConnectorload_from_checkpoint(start, end, ckpt)增量 + 断点续传大源头的主力路径

定义分别在 backend/onyx/connectors/interfaces.py:120LoadConnector)、:125PollConnector)、:134SlimConnector)、:266CheckpointedConnector)。

两个"带权限同步"的变体是给企业版走 ACL 的:SlimConnectorWithPermSync.retrieve_all_slim_docs_perm_sync:147)在拿 ID 的同时把每个文档的外部权限带回来;CheckpointedConnectorWithPermSync.load_from_checkpoint_with_perm_sync:305)则在正常抓取时顺带把权限一起灌进 Document

ConnectorRunner 会检查这一点:构造时如果 include_permissions=True 但连接器不是 CheckpointedConnector,直接抛 ValueErrorbackend/onyx/connectors/connector_runner.py:117)。

3.2 辅助契约:不是"怎么抓",而是"抓之前/之外的事"

契约解决什么位置
OAuthConnector提供授权 URL 和 code→token 换取,让前端能跑 OAuth 流程interfaces.py:158
CredentialsConnector凭据会在跑的过程中轮换(如短期 token)interfaces.py:240
EventConnector监听推送事件而不是轮询interfaces.py:253
HierarchyConnector单独吐"目录树"节点(空间 / 文件夹 / 频道)interfaces.py:332
Resolver针对一批已知失败的文档重抓,不带 checkpointinterfaces.py:316

EventConnector预留的——backend/onyx/connectors/README.md 里明说后台任务目前不用它。这是诚实的边界,不要以为 Onyx 已经支持 webhook 驱动索引。

Resolver 是 §7.4 定向重建的入口:它接收一组 ConnectorFailure,尽力把这些文档重新吐出来,由调用方负责把旧的失败记录换成新的。

3.3 基类里的公共动作

BaseConnectorinterfaces.py:43)不只是空壳,它塞了几个所有连接器共享的钩子:

  • parse_metadata:53)把元数据字典拍平成 key: value 行,非字符串就直接抛错——逼连接器自己实现解析。
  • validate_perm_sync:79)显式写着"别覆盖这个方法",它内部用 fetch_ee_implementation_or_noop 去企业版包里找实现,社区版就是空操作。这是 Onyx 处理开源/企业分层的通用招式。
  • normalize_url:105)返回 NormalizationResult(use_default=True) 表示"我没实现,用默认归一化器"——用返回值而不是异常来表达"未实现",调用方不用 try/except。
  • set_raw_file_callback:95)注入一个"把原始字节存下来"的回调,不关心的连接器就让它保持 None

3.4 注册与实例化:懒加载 + 按输入类型校验

注册表是纯数据CONNECTOR_CLASS_MAP 里每一项只是 ConnectorMapping(module_path=..., class_name=...) 两个字符串(backend/onyx/connectors/registry.py:8)。

工厂负责真正 import_load_connector_classbackend/onyx/connectors/factory.py:40)用 importlib.import_module 动态加载并缓存到 _connector_cache。好处很实在——不装 Salesforce SDK 的部署,也不会因为 import 失败而起不来。

实例化前还会校验"这个连接器支不支持你要的输入类型",_validate_connector_supports_input_typefactory.py:59):

# 示意,非源码:poll 的校验为什么要放两个条件
poll_unsupported = (
input_type == InputType.POLL
and not issubclass(connector, PollConnector) # 老式增量
and not issubclass(connector, CheckpointedConnector) # 新式增量
)

重点看:CheckpointedConnector 被当成 POLL 的合法实现——源码里那行注释直说了"将来所有连接器都该是 checkpoint 连接器",这是一次进行中的迁移

3.5 动态凭据:用 Redis 锁围住 token 轮换

instantiate_connectorfactory.py:106)分两条路:

  • 实现了 CredentialsConnector → 注入一个 OnyxDBCredentialsProvider,连接器自己按需读写凭据。
  • 否则 → 直接解密一次 credential_json 塞给 load_credentials,如果连接器返回了新凭据就写回数据库。

OnyxDBCredentialsProviderbackend/onyx/connectors/credentials_provider.py:17)的关键是那把锁:

self.lock_key = f"da_lock:connector:{connector_name}:credential_{credential_id}"
self._lock: RedisLock = self.redis_client.lock(self.lock_key, self.LOCK_TTL)

credentials_provider.py:34-35LOCK_TTL = 900 秒。它解决的是:同一份 refresh token 被两个并发的索引任务同时拿去换新 token,后换的那个会让先换的失效。is_dynamic() 返回 True 就是在告诉调用方"你必须用锁"(接口文档见 interfaces.py:229)。静态凭据走 OnyxStaticCredentialsProvider:106),is_dynamic() 返回 False,锁退化成空操作。


4. 断点续传:generator 的返回值就是存档点

它要解决的小问题: 一个连接器跑 3 小时,跑到第 2 小时被 OOM 杀掉。下次重启,怎么接着跑而不是从头再来?

4.1 思路:用 Python generator 的 return 值当 checkpoint

CheckpointedConnector.load_from_checkpoint 的类型是 Generator[Document | HierarchyNode | ConnectorFailure, None, CT]interfaces.py:259CheckpointOutput 别名)。三段式的含义:

  • yield 出去的:文档、层级节点、或者一条失败记录。
  • 最后 return 的:新的 checkpoint 对象。

这个设计的妙处在于类型系统帮你保证了"checkpoint 有且只有一个,且一定在最后"——连接器作者没法在中间偷偷 yield 一个 checkpoint。源码注释里就是这么解释的(interfaces.py:274-292)。

代价是:Python 里拿 generator 的返回值很别扭,要么 yield from,要么捕 StopIteration.value。所以有了包装器。

4.2 CheckpointOutputWrapper:把三合一的流拆成四元组

backend/onyx/connectors/connector_runner.py:54CheckpointOutputWrapper 把上面那个 generator 转成统一的四元组流:

连接器 yield 的东西 包装后 yield 的四元组
─────────────────────────────────────────────────────
Document ──▶ (doc, None, None, None)
HierarchyNode ──▶ (None, node, None, None)
ConnectorFailure ──▶ (None, None, fail, None)
(最后的 return) ──▶ (None, None, None, checkpoint)

内部用一个 _inner_wrapperself.next_checkpoint = yield from ... 把返回值截下来(:74-78),流结束后再单独 yield 一次。如果连接器忘了 return checkpoint,这里直接 RuntimeError:90-93)——宁可炸也不要静默地丢掉存档点

4.3 ConnectorRunner.run:四种连接器一个出口

ConnectorRunner.runconnector_runner.py:130)是"把四类契约折叠成一个接口"的地方。它做三件事:成批、补日志、统一出口形状(类注释 :99-104)。

对非 checkpoint 连接器,它伪造一个已完成的存档点:

finished_checkpoint = self.connector.build_dummy_checkpoint()
finished_checkpoint.has_more = False

connector_runner.py:229-230。之后 PollConnectorLoadConnector 各跑一遍自己的 generator,末尾 yield 这个假 checkpoint(:240:249)。上层的 while checkpoint.has_more 循环因此对四种连接器一视同仁。

批次里还藏着一条父先于子的不变量:层级节点批次总是在文档批次之前 flush,因为文档要引用父节点 ID:

# 攒够一批文档时,先把手里的层级节点冲出去,确保父节点已存在
if len(self.doc_batch) >= self.batch_size:
if len(self.hierarchy_node_batch) > 0:
yield None, self.hierarchy_node_batch, None, None
self.hierarchy_node_batch = []
yield self.doc_batch, None, None, None

connector_runner.py:204-209_separate_batch:276)则负责把老式连接器吐出来的混合 list 拆成文档和节点两堆。

4.4 失败降级:一条记录,而不是一次崩溃

这是整章最值得学的一条:失败是数据,不是控制流

ConnectorFailurebackend/onyx/connectors/models.py:535)有个校验器强制"要么指明失败的文档、要么指明失败的实体,二者必居其一且只能其一"(:548)。它带一个被排除出序列化的 exception 字段——能上报 Sentry,但不会污染 JSON。

抓取循环拿到 failure 时的动作(run_docfetching.py:710-742):打 Sentry tag、create_index_attempt_error 落库、计数加一,然后继续循环

ConnectorRunner.run 那个大 exceptconnector_runner.py:258-280)反而不吞异常——它做的是取最深一层 traceback 的局部变量打进日志(截断到 1024 字符),然后 raise。区别很清楚:

  • 连接器主动 yield 的失败 → 降级成记录,继续。
  • 连接器抛出的异常 → 一路上抛,整个 attempt 失败。

run_docfetching.py:993-1005 的注释解释了为什么不能对后者也降级:黑盒 raise 里没法定位"是哪个实体坏了",把 attempt 悄悄标成 COMPLETED_WITH_ERRORS 会让系统误以为源数据已经抓全。

4.5 checkpoint 存的是"发出去了",不是"索引好了"

这是很容易读错的一处,源码专门加了注释(run_docfetching.py:868-870):

checkpoint 追踪的是哪些批次已经送进 filestore,不是哪些批次已经索引完成

配套逻辑在 :563-600:恢复时如果确实是从 checkpoint 续跑,就调 reissue_old_batches 把之前存好但没处理完的批次重新派发一遍,并把 last_batch_num 接上去。非 checkpoint 连接器则直接 cleanup_all_batches()——反正重跑会把那些文档再抓一遍。

时间窗还有个反直觉的细节:上次失败的话,这次沿用上次的时间窗:513-523)。原因写在注释里——Slack 新频道这类信息是缓存在 checkpoint 里的,换了窗口就会漏掉。


5. 索引管线主干:一批文档的十步

5.1 三层入口

run_indexing_pipeline(...) ← 读 search_settings,决定用哪套嵌入/是否开 contextual RAG


index_doc_batch_with_handler(...) ← 唯一职责:把异常包成 ConnectorFailure


index_doc_batch(...) ← 真正干活

run_indexing_pipelinebackend/onyx/indexing/indexing_pipeline.py:1532)做的都是"选配置":主索引还是次索引(模型切换期间会有 FUTURE 状态的次索引,:1492-1498)、是否 multipass、contextual RAG 用哪个 LLM,最后现造一个 Chunker

index_doc_batch_with_handler:387)只有一个 try/except,但语义很关键:ConnectorStopSignal 原样上抛(用户按了暂停),其它任何异常都被转成这一批每篇文档各一条 ConnectorFailure:431-448)。整批失败也是失败记录,不是崩溃。

5.2 index_doc_batch 的十步

怎么读这张图:从上到下,左边是步骤,右边是"这一步可能把哪些文档踢出去"。

① filter_documents → 空文档丢弃、超长文档 → Failure
② 文档摄入 hook → 外部服务可改写/否决文档
③ adapter.prepare → 去重 + 写 Postgres,全都跳过就提前返回
④ process_image_sections → 图片走视觉 LLM 变成文字
⑤ chunker.chunk → 切成 DocAwareChunk
⑥ add_contextual_summaries→ 用 LLM 给每个 chunk 补文档级上下文(可选)
⑦ embed_and_stream → 嵌入,逐批落盘;失败的整篇文档被剔除
⑧ adapter.lock_context → 行锁,防止并发改同一批文档
⑨ write_chunks_to_vector_db_with_backoff → 写库;失败降级到逐文档重试
⑩ post_index + 回写 content_hash → 只给确认写成功的文档盖章

主体在 indexing_pipeline.py:1218-1470

5.3 增量判定:get_docs_to_update 的两道闸门

它要解决的小问题: 每次同步都重新嵌入所有文档太贵。怎么判断"这篇没变"?

get_docs_to_updateindexing_pipeline.py:320)用两道闸门,顺序有讲究:

闸门依据什么时候生效代价
闸门 1源头给的 doc_updated_at 没往前走连接器提供了时间戳时零成本,不用读正文
闸门 2content_hash() 和库里存的一样时间戳缺失或没前进时要算一次 MD5

真正精妙的是闸门 2 什么时候不生效

# 时间戳前进 = 权威证据,此时绝不用 hash 去推翻它
content_hash = doc.content_hash()
if not timestamp_advanced:
db_doc = id_to_db_doc_map.get(doc.id)
if db_doc and db_doc.content_hash == content_hash:
continue # 内容没变,跳过

indexing_pipeline.py:374-379。docstring(:340-343)给了具体反例:Google Drive 原地替换图片,image_file_id 不变、图片字节变了——hash 会错误地说"跳过",但时间戳前进了,所以必须以时间戳为准。

content_hash 的定义在 backend/onyx/connectors/models.py:297。它故意不包含图片摘要——摘要是 LLM 生成的、不确定的,而且算 hash 的时候还没生成。

doc_id_to_content_hash 算出来后不立刻写库,而是一路带到第 ⑩ 步,只给"确认写进向量库的文档"盖章(indexing_pipeline.py:1438-1449)。注释说得很明白:提前存 hash 会让一次失败的索引永久地跳过这篇文档。

5.4 丢弃规则:filter_documents

filter_documentsindexing_pipeline.py:590)区分了"静默丢弃"和"记为失败":

情况处理为什么
既无标题也无正文静默丢弃没有任何有用信息
标题是显式空串且正文空静默丢弃拼出来的 chunk 文本会是空的
总字符数 > MAX_DOCUMENT_CHARS记一条 ConnectorFailure用户需要知道并去拆分文档

MAX_DOCUMENT_CHARS 默认 536,870,912(512MB 字符,backend/onyx/configs/app_configs.py:1510)。源码注释解释了为什么这个兜底必须在这里而不只在连接器里:后面几步内存开销大,一篇超长文档能把整个容器 OOM 掉(:641-644)。失败消息里直接写了"限制由 MAX_DOCUMENT_CHARS 设定"和"把文档拆小",是给管理员看的可执行提示。

5.5 图片:视觉 LLM 把像素变成可检索文字

process_image_sectionsindexing_pipeline.py:696)的产物是 IndexingDocument——比 Document 多一个 processed_sections 字段,原始 section 保留、处理后的 section 另放。

三条路径:

  • 没开图片分析 / 批次里没图片 → 不去拿视觉 LLM,图片 section 变成空文本的基础 Section:717-750)。
  • 开了但没有视觉模型 → 打一条 warning 告诉管理员去配(:724-729),仍然降级。
  • 正常 → 每张图片读出字节,攒进 pending,最后用 run_functions_tuples_in_parallel(..., allow_failures=True, max_workers=MAX_IMAGE_WORKERS) 并行摘要(:816-820MAX_IMAGE_WORKERS = 16)。

有个容易忽略的正确性细节:判断是不是图片用的是 section.type == SectionType.IMAGE 而不是 isinstance,因为 section 经过 pydantic 往返后会变成基类实例(注释在 :710-712)。同理,非图片 section 用 section.model_copy() 而不是重建基础 Section——否则会丢掉 TabularSectioncsv_file_id:765-766)。

失败也不中断:读不到文件就把文本设成 "[Image could not be processed]",异常就设成 "[Error processing image]"——这些占位符会照样被索引。

5.6 Contextual RAG:让每个 chunk 知道自己在讲什么

它要解决的小问题: 一个 chunk 里写着"该指标下降了 12%",但"该指标"是什么、这是哪份报告,都在别的 chunk 里。单看这个 chunk,检索命中率很差。

思路: 索引时就用 LLM 给每个 chunk 补一段"我在这篇文档里扮演什么角色"的说明,拼进被嵌入的文本。

三个函数分工(都在 indexing_pipeline.py):

函数产出每篇文档调用 LLM 几次
add_document_summaries:828整篇文档摘要,写进每个 chunk 的 doc_summary1 次
add_chunk_summaries:870每个 chunk 的 chunk_contextchunk 数量次(并行)
add_contextual_summaries:953编排:按文档分组、算 token 预算、按开关调上面两个

add_chunk_summaries 里有两个值得学的细节。

一是省 token 的降级链:893-911):文档 token 数 ≤ MAX_TOKENS_FOR_FULL_INCLUSION(4096)就直接塞全文;超了就用上一步算好的摘要;如果文档摘要开关是关的、又拿不到摘要,才临时调一次 LLM 算摘要。

二是 prompt 缓存的切分:文档上下文是 CONTEXTUAL_RAG_PROMPT1(对同一篇文档的所有 chunk 都一样),chunk 内容是 CONTEXTUAL_RAG_PROMPT2。前者作为 cacheable_prefix 传给 process_with_prompt_cache,后者作为 suffix(:922-927)。一篇 200 个 chunk 的文档,前缀只付一次全价。同一套 prompt 缓存思路在对话侧也用,见 02 章

失败姿态assign_context 里捕住 LLMRateLimitError 和一般异常,把 chunk_context 设成空串继续(:938-945),注释直说"在 chunker 阶段报错是不可接受的"。宁可少一点上下文,也不能让整批索引挂掉。

并发上限 MAX_CONTEXTUAL_RAG_WORKERS = 128,旁边注释写着"假设每个 worker 8MB 内存"(:102)——这个数是按内存反推的。

5.7 完备性断言:不许有文档"人间蒸发"

_verify_indexing_completenessindexing_pipeline.py:995)在每个索引写完后做一次集合等式检查:

{成功插入的 doc_id} ∪ {写库失败的} ∪ {嵌入失败的} 必须 == {本应更新的 doc_id}

不等就 RuntimeError,错误消息里带着"This should never happen"。这是一条内部一致性断言:每篇文档要么成功、要么有一条明确的失败记录,不允许静默消失。

5.8 两个外部钩子

管线上留了两个给外部服务插手的点:

钩子时机能做什么
DOCUMENT_INGESTION入库之前改写 section,或者返回空 section 直接否决这篇文档
DOCUMENT_PUSH成功写库之后把文档推给外部目的地

_apply_document_ingestion_hook:1018)有个省钱技巧:先只跑第一篇文档,如果返回 HookSkipped(说明钩子压根没配),立刻返回,省掉剩下 N-1 次数据库查询(:1115-1126)。

_maybe_push_documents:1148)的三道闸门都是安全考虑(:1160-1178):初次索引不推(from_beginning)、多租户环境不推(会把不同组织的数据混进同一个外部目的地)、非 public 的 cc_pair 不推。


6. 分块与嵌入

6.1 token 预算:chunk 里到底能放多少正文

它要解决的小问题: 嵌入模型的输入窗口是固定的 512 token(DOC_EMBEDDING_CONTEXT_SIZEbackend/shared_configs/configs.py:42)。标题要占、元数据要占、contextual RAG 的摘要还要占——正文剩多少?

Chunker._handle_single_documentbackend/onyx/indexing/chunker.py:193)是一段逐级让步的预算算术。怎么读这张图:从上到下依次让步,任何一步发现正文空间不够就退回上一层的方案。

起点:content_token_limit = 512

① 减去标题 token

② 减去元数据 token
│ ├─ 元数据 ≥ 512 × 0.25 ? → 元数据整段丢弃,不进语义文本

③ 减去 contextual RAG 预留
│ └─ 整篇文档能装进一个 chunk ? → 不预留(没必要补上下文)

④ 剩下的正文空间 ≤ 256 (CHUNK_MIN_CONTENT) ?
│ └─ 是 → 放弃 contextual RAG,退回 ② 的预算

⑤ 还是 ≤ 256 ?
└─ 是 → 标题前缀和元数据后缀全部清空,正文独占 512

对应源码:MAX_METADATA_PERCENTAGE = 0.25CHUNK_MIN_CONTENT = 256chunker.py:30-36;元数据超限丢弃在 :218-220;单 chunk 装得下就不预留在 :229-233;两次 <= CHUNK_MIN_CONTENT 的让步在 :251-263

一句话总结这段设计:正文永远优先。宁可丢掉标题和元数据这些辅助信号,也要保证 chunk 里有足够多的真实内容。

6.2 元数据怎么变成可检索文本

get_metadata_suffix_for_document_indexchunker.py:41)同时产出两个串:

  • 语义串"Metadata:\n\tkey - value\n...",人话格式,拼进被嵌入的文本。
  • 关键词串:只有值、空格分隔,给关键词检索用。

分开的理由很直接:向量模型需要看到"key 是什么"才能理解语义;BM25 只关心值里的词,键名反而是噪音。

最终被送去嵌入的完整文本由 generate_enriched_content_for_chunk_embedding 拼出(backend/onyx/document_index/chunk_content_enrichment.py:11):

title_prefix + doc_summary + content + chunk_context + metadata_suffix_semantic

这是写入侧和读取侧的接缝。 同一个文件里的 cleanup_content_for_chunks:17)是它的逆运算——检索返回给用户之前,要把这些索引期添加物按相反顺序剥掉,否则用户看到的引用里会混着"Metadata:"和 LLM 生成的摘要。读取侧怎么用这些字段见 04 章

6.3 大 chunk 与 mini chunk:同一份内容的三种粒度

粒度怎么来的干什么用
mini chunkmini_chunk_splitter 按 150 token 再切提高细粒度匹配召回
普通 chunkchunk_splitter 按 512 token 切主力
large chunk连续 4 个普通 chunk 合并捕获跨 chunk 的语义

generate_large_chunkschunker.py:111)按 LARGE_CHUNK_RATIO = 4 分组,只有组内多于 1 个 chunk 才合并:114)——避免造出一个和原 chunk 一模一样的"大" chunk。

_combine_chunks:70)里最容易写错的是链接偏移量:合并后每个原 chunk 的 source_links 偏移都要加上前面所有内容的长度:

offset += len(SECTION_SEPARATOR) + len(chunks[i - 1].content)
for link_offset, link_text in (chunks[i].source_links or {}).items():
merged_chunk.source_links[link_offset + offset] = link_text

chunker.py:102-106。偏移算错了,引用就会指到文本里的错误位置。

大 chunk 和 mini chunk 互斥:embed_chunks 里如果发现一个大 chunk 带着 mini chunk,直接 RuntimeError("Large chunk contains mini chunks")backend/onyx/indexing/embedder.py:149)。注释解释了为什么这不是损失——大 chunk 高分命中时本来就不会用 mini chunk。

6.4 嵌入:本进程只负责拼请求

DefaultIndexingEmbedderembedder.py:88不加载任何模型。它构造的 EmbeddingModel 明确指向索引专用的 model server:

server_host=INDEXING_MODEL_SERVER_HOST,
server_port=INDEXING_MODEL_SERVER_PORT,
retrim_content=True,

embedder.py:72-74。真正的分工是:

docprocessing worker model server(独立进程/容器)
┌──────────────────────┐ ┌─────────────────────────┐
│ 拼 chunk 文本 │ EmbedRequest │ 本地模型推理 │
│ 缓存标题向量 │ ────────────▶ │ 或 │
│ 把向量映射回 chunk │ ◀──────────── │ 转发给 OpenAI/Cohere... │
└──────────────────────┘ EmbedResponse└─────────────────────────┘

请求组装在 backend/onyx/natural_language_processing/search_nlp_models.py:954_batch_encode_texts,多线程发批;encode:1130)在发送前做三件清理:大 chunk 存在时把 max_seq_length 乘以 LARGE_CHUNK_RATIO、按 tokenizer 裁剪超长文本、剔除会让 UTF-8 编码炸掉的非法 Unicode(:1160-1161)。

标题向量只算一次embed_chunks 先把批次里所有不重复的标题收成集合,一次性嵌入并缓存进 title_embed_dictembedder.py:160-185)。一篇 200 chunk 的文档,标题就省了 199 次调用。如果走到"标题得单独嵌入"的兜底分支,会打 logger.error("...this should not happen!")——把不该发生的事显式暴露出来。

from_db_search_settings:228)是唯一的正规构造入口:嵌入配置全部来自数据库里的 SearchSettings,这样"切换嵌入模型"才可能是一次配置变更(见 §7.4)。

6.5 三层失败隔离 + 落盘

嵌入这一段有三层保护,逐层收窄:

第 1 层:整批嵌入
└─ 失败 → sleep 2 秒,改成按文档逐篇嵌入
第 2 层:逐篇嵌入
└─ 某篇失败 → 记一条 ConnectorFailure,其它篇继续
第 3 层:跨批次清理
└─ 第 N 批失败的文档,第 1..N-1 批里它的 chunk 也要被抹掉

第 1、2 层在 embed_chunks_with_failure_handlingembedder.py:251);第 3 层在 _embed_chunks_to_storeindexing_pipeline.py:209)。

第 3 层的必要性来自每篇文档必须全有或全无all_failed_doc_ids 跨批次累积,后来的批次直接跳过这些文档(:232-235),已经落盘的早期批次则由 store.scrub_failed_docs 回头擦掉(:265-266)。半篇文档写进索引,比整篇缺失更糟——用户会搜到残缺内容还以为是全部。

为什么要落盘? ChunkBatchStorebackend/onyx/indexing/chunk_batch_store.py:10)把每批嵌入好的 chunk pickle 到临时目录。向量是几百上千维的浮点数组,一批 1000 个 chunk(MAX_CHUNKS_PER_DOC_BATCH)全留在内存里很容易 OOM。stream():66)每次调用返回一个新的 generator,所以在模型切换期间可以为主索引和次索引各遍历一遍。目录生命周期由 embed_and_stream 这个 contextmanager 兜住(indexing_pipeline.py:279),退出即删。

写向量库那一层用同样的招式:write_chunks_to_vector_db_with_backoffbackend/onyx/indexing/vector_db_insertion.py:31)先试整批,失败就 sleep 2 秒后按 groupby(doc_id) 逐文档重试。它还有一个顺序断言——同一个 doc_id 出现两次就 RuntimeError:85-88),因为 groupby 只对连续分组有效,chunk 顺序错乱会静默漏数据。507 Insufficient Storage 有专门的提示日志,直接告诉运维"给索引容器加内存或磁盘"(:22-30)。


7. 生命周期的另一半:删、剪、改权限、重建

前面六节讲的都是"写进去"。索引要长期正确,还得有四条把索引拉回和源头一致的路径。

7.1 剪枝:源头删了的文档,怎么从索引里去掉

思路: 把源头当前所有 ID 拉一遍,和本地记录做差集。

connector_pruning_generator_taskbackend/onyx/background/celery/tasks/pruning/tasks.py:481)用 InputType.SLIM_RETRIEVAL 实例化连接器,然后:

doc_ids_to_remove = list(all_indexed_document_ids - all_connector_doc_ids.keys())

pruning/tasks.py:690-692

extract_ids_from_runnable_connectorbackend/onyx/background/celery/celery_utils.py:146)负责"用现有契约凑出一份全量 ID 列表",优先级依次是 SlimConnectorSlimConnectorWithPermSyncLoadConnectorPollConnector(时间窗设成 1970 到现在)→ CheckpointedConnector,都不是就 RuntimeError

有一个安全默认值值得单独拎出来:抓 ID 的过程中遇到 ConnectorFailure,它把失败文档的 ID 也塞进 ID 集合(celery_utils.py:118-124)。因为"抓失败"和"确实不存在"在这里长得一样,如果不这么做,一次限流就会误删一批文档。

会话管理也很讲究:连接器爬 10-30 分钟期间不持有数据库连接——枚举前一个 session 拿配置,枚举后再开一个 session 做差集和派发(注释在 pruning/tasks.py:560-563)。

7.2 删除:一个文档被几个连接器共享时怎么办

document_by_cc_pair_cleanup_taskbackend/onyx/background/celery/tasks/shared/tasks.py:111)按引用计数分岔:

该文档被几个 cc_pair 索引动作
恰好 1 个从所有索引里删除,再删 Postgres 记录
多于 1 个只更新元数据(去掉这个 cc_pair 带来的 ACL),文档留着

任务被切成三段,每段之间释放数据库连接(注释在 :110-116):先读 DB 状态,再做向量库 HTTP 调用,最后开新事务写回。原因写得很直白——跨 HTTP 往返持有 pg 事务会钉住一个 pgbouncer 槽位,批量删除时整个连接池会被打满。

7.3 权限:ACL 是怎么变成 chunk 上的字段的

ACL 字符串在写入时才拼出来。 DocumentAccess.to_acl()backend/onyx/access/models.py:176)把五类来源统一成一个带前缀的字符串集合:

来源前缀函数
Onyx 用户邮箱prefix_user_email
Onyx 用户组prefix_user_group
外部系统用户邮箱prefix_user_email
外部系统组 IDprefix_external_group
公开常量 PUBLIC_DOC_PAT

方法的 docstring 明确要求:查询时构造的 ACL 过滤串必须用同样的格式——这是写入侧和读取侧之间一条不成文但致命的约定。

真正把 ACL 挂到 chunk 上的是适配器。DocumentIndexingBatchAdapter.prepare_enrichmentbackend/onyx/indexing/adapters/document_indexing_adapter.py:115)一次性把整批的权限、文档集、层级祖先、boost、旧 chunk 数全查出来,做成一个 DocumentChunkEnricher;后者的 enrich_chunk:272)纯内存地把这些拼进 DocMetadataAwareIndexChunkbackend/onyx/indexing/models.py:94)。N 个 chunk 只做 1 轮数据库查询。

默认值是 no_accessdocument_indexing_adapter.py:133-139):查不到权限信息的文档,is_public=False 且没有任何用户/组——默认谁都看不见,而不是默认公开。

独立的权限同步路径:企业版的 element_update_permissionsbackend/ee/onyx/background/celery/tasks/doc_permission_syncing/tasks.py:689)负责在不重新索引的前提下刷新权限。它先把外部用户批量落库(continue_on_error=True),再 upsert 文档的外部权限;如果这次 upsert 创建了新文档记录,还会顺手把它和 cc_pair 关联起来(:709-717)——处理的是"权限同步先于内容同步发现了一个新文档"的情况。

source_should_fetch_permissions_during_indexingbackend/onyx/access/access.py:145)决定某个源是"索引时顺带抓权限"还是"交给独立的权限同步任务"。抓取侧还叠了一条策略:已经成功索引过的 cc_pair 就不在索引时抓权限了,交给 doc_sync(run_docfetching.py:526-533)。

7.4 换模型与重建:两条不同的路

整体换嵌入模型 走的是"双索引并行":run_indexing_pipeline:1492-1498 检查次索引状态是不是 FUTURE,是就用次索引的 search settings;index_doc_batch 则对 document_indices 列表里的每个索引都写一遍(:1388)。新旧索引同时被灌数据,切换时才翻转。

Vespa → OpenSearch 的存量迁移 是一个独立的周期任务 migrate_chunks_from_vespa_to_opensearch_taskbackend/onyx/background/celery/tasks/opensearch_migration/tasks.py:73)。它不重新嵌入——用 Vespa 的 Visit API 把 chunk 整块搬过去,transform_vespa_chunks_to_opensearch_chunksbackend/onyx/background/celery/tasks/opensearch_migration/transformer.py:184)做字段格式转换(包括 ACL,见 _transform_vespa_acl_to_opensearch_acl:172)。进度按分片存 continuation token,状态表操作在 backend/onyx/db/opensearch_migration.py

定向重建是最细粒度的一条:process_targets_for_cc_pairbackend/onyx/background/indexing/run_targeted_reindex.py:214)拿一组具体的失败文档 ID,把它们转成 ConnectorFailure 喂给 Resolver.reindex,再走标准管线。连接器没实现 Resolver 就直接把这些目标标成 still_failing,让管理员清楚看到"这个源还不支持定向重建"(:260-266),而不是静默无事发生。


8. 巧妙之处(可以带走的技术)

① 用 generator 的返回值当 checkpoint。 类型 Generator[Doc|Node|Failure, None, CT] 在编译期就保证了"checkpoint 有且只有一个、且在最后",连接器作者无法写错。代价是取返回值别扭,于是配一个 CheckpointOutputWrapper 消化掉(connector_runner.py:54)。

② 失败是数据,不是控制流。 ConnectorFailure 带一个 exclude=True 的 exception 字段(connectors/models.py:538),能上报 Sentry 又不污染序列化;model_validator 强制"文档失败"和"实体失败"二选一(:548)。整条管线的默认反应是"记一笔,继续"。

③ 时间戳压过内容哈希。 两道去重闸门里,时间戳前进被当作权威证据,主动禁用哈希闸门(indexing_pipeline.py:374-379)。这是一条从真实 bug(GDrive 原地换图)里长出来的规则。

④ hash 只给确认写成功的文档盖章。 提前存 hash = 一次失败就永久跳过(indexing_pipeline.py:1434-1449)。凡是"用来判断能不能跳过"的状态,都该在动作成功之后才落库。

⑤ 整批 → 逐条的两级退避,用了三次。 嵌入(embedder.py:265-294)、写向量库(vector_db_insertion.py:45-99)、图片摘要(allow_failures=True)都是同一个形状:快路径整批走,出错才付逐条的代价来定位问题。

⑥ 每篇文档全有或全无。 跨批次的 scrub_failed_docschunk_batch_store.py:75)保证不会出现"半篇文档进了索引"。

⑦ 嵌入结果落盘而不是留在内存。 用磁盘换内存上限,还顺带让"同一份向量写两个索引"变得免费(chunk_batch_store.py:66stream 每次返回新 generator)。

⑧ 钩子先探一次再全量跑。 第一篇文档返回 HookSkipped 就直接收工,省掉 N-1 次数据库查询(indexing_pipeline.py:1115-1126)。

⑨ 副作用被赶进适配器。 IndexingBatchAdapterbackend/onyx/indexing/models.py:255)把"写 Postgres、加锁、拼 ACL、回写计数"抽成协议,同一条管线因此能同时服务连接器文档和用户上传文件两种场景,各自有一个实现(backend/onyx/indexing/adapters/)。

⑩ 每个阶段自己开短会话。 适配器每个方法开关一次自己的 session(类注释 document_indexing_adapter.py:48-49),嵌入和写库那段长耗时完全不持有数据库连接。


9. 边界与局限(诚实版)

  • EventConnector 是空占位。 README.md 明说后台任务不用它,interfaces.py:253 只有一个抽象方法。Onyx 目前没有 webhook 驱动的索引。
  • checkpoint 只能是内存里的 pydantic 对象。 ConnectorCheckpoint 上挂着 TODO:"也许该改成基于磁盘的,以应对超大 checkpoint"(connectors/models.py:511-512)。代码里有 check_checkpoint_size 定期检查,但那是报警不是解决。
  • contextual RAG 的重试很弱。 add_chunk_summaries 遇到限流就把上下文设成空串,# TODO: for v2, add robust retry logicindexing_pipeline.py:940)。开着这个功能的大规模索引会静默地降质。
  • 限流靠字符串匹配识别。 celery_utils.py:203-208 的注释坦白:连接器报限流的方式五花八门,现在靠在错误消息里找 "rate limit" 或 "429",TODO 是引入统一的 ConnectorRateLimitError
  • 失败阈值那段代码可能已经没用了。 run_docfetching.py:328 顶上写着 # TODO: delete from here if ends up unused,而且开了 PERSISTENT_INDEXING 就直接 return。
  • section_continuation 已废弃。 backend/onyx/indexing/models.py:40-42 标了"OpenSearch 迁移后已废弃,不要用",但字段还在。
  • chunk 文本清理是脆的。 cleanup_content_for_chunks 靠字符串前后缀匹配来剥掉索引期添加物,作者自己写了"这整个函数都不怎么样,推 OpenSearch 前要在 QA 阶段清理掉"(chunk_content_enrichment.py:40-41)。
  • _verify_indexing_completeness 抛的是硬错误。 集合对不上就 RuntimeError,整批失败。它是一致性守卫,不是恢复机制。

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

主题文件路径符号名
连接器基类与公共钩子backend/onyx/connectors/interfaces.pyBaseConnectorvalidate_perm_syncnormalize_url
四类主流契约backend/onyx/connectors/interfaces.pyLoadConnectorPollConnectorSlimConnectorCheckpointedConnector
权限同步变体backend/onyx/connectors/interfaces.pySlimConnectorWithPermSyncCheckpointedConnectorWithPermSync
辅助契约backend/onyx/connectors/interfaces.pyOAuthConnectorCredentialsConnectorEventConnectorHierarchyConnectorResolver
checkpoint 类型别名backend/onyx/connectors/interfaces.pyCheckpointOutput
数据源 → 类的映射表backend/onyx/connectors/registry.pyCONNECTOR_CLASS_MAPConnectorMapping
懒加载与实例化backend/onyx/connectors/factory.py_load_connector_classinstantiate_connector_validate_connector_supports_input_type
动态凭据与 Redis 锁backend/onyx/connectors/credentials_provider.pyOnyxDBCredentialsProviderOnyxStaticCredentialsProvider
四种连接器折叠成一个出口backend/onyx/connectors/connector_runner.pyConnectorRunner.run_separate_batch
拆解 checkpoint generatorbackend/onyx/connectors/connector_runner.pyCheckpointOutputWrapperbatched_doc_ids
失败与存档点的数据模型backend/onyx/connectors/models.pyConnectorFailureConnectorCheckpointDocumentBase.content_hash
抓取主循环backend/onyx/background/indexing/run_docfetching.pyconnector_document_extraction_get_connector_runner
加工任务入口backend/onyx/background/celery/tasks/docprocessing/tasks.pydocprocessing_task_docprocessing_task
管线三层入口backend/onyx/indexing/indexing_pipeline.pyrun_indexing_pipelineindex_doc_batch_with_handlerindex_doc_batch
增量判定backend/onyx/indexing/indexing_pipeline.pyget_docs_to_updateindex_doc_batch_prepare
丢弃规则backend/onyx/indexing/indexing_pipeline.pyfilter_documents
图片摘要backend/onyx/indexing/indexing_pipeline.pyprocess_image_sections
contextual RAGbackend/onyx/indexing/indexing_pipeline.pyadd_contextual_summariesadd_document_summariesadd_chunk_summaries
一致性守卫与外部钩子backend/onyx/indexing/indexing_pipeline.py_verify_indexing_completeness_apply_document_ingestion_hook_maybe_push_documents
分块与 token 预算backend/onyx/indexing/chunker.pyChunker._handle_single_document_get_metadata_suffix_for_document_index
大 chunk 合并backend/onyx/indexing/chunker.pygenerate_large_chunks_combine_chunks
section → chunk 分派backend/onyx/indexing/chunking/document_chunker.pyDocumentChunker.chunk_select_chunker
嵌入与失败隔离backend/onyx/indexing/embedder.pyDefaultIndexingEmbedder.embed_chunksfrom_db_search_settingsembed_chunks_with_failure_handling
model server 请求组装backend/onyx/natural_language_processing/search_nlp_models.pyEmbeddingModel.encode_batch_encode_texts
向量落盘中转backend/onyx/indexing/chunk_batch_store.pyChunkBatchStore.streamscrub_failed_docs
写向量库退避backend/onyx/indexing/vector_db_insertion.pywrite_chunks_to_vector_db_with_backoff
心跳与停止信号backend/onyx/indexing/indexing_heartbeat.pyIndexingHeartbeatInterface
兜底失败记录backend/onyx/indexing/persistent_indexing.pybuild_generic_connector_failurerecord_generic_failure
副作用适配器backend/onyx/indexing/adapters/document_indexing_adapter.pyDocumentIndexingBatchAdapterDocumentChunkEnricher.enrich_chunk
适配器协议backend/onyx/indexing/models.pyIndexingBatchAdapterDocMetadataAwareIndexChunk
ACL 字符串生成backend/onyx/access/models.pyDocumentAccess.to_aclExternalAccess
索引期文本拼装 / 剥离backend/onyx/document_index/chunk_content_enrichment.pygenerate_enriched_content_for_chunk_embeddingcleanup_content_for_chunks
剪枝backend/onyx/background/celery/tasks/pruning/tasks.pyconnector_pruning_generator_task
全量 ID 枚举backend/onyx/background/celery/celery_utils.pyextract_ids_from_runnable_connector
删除 / 引用计数backend/onyx/background/celery/tasks/shared/tasks.pydocument_by_cc_pair_cleanup_task
权限写回backend/ee/onyx/background/celery/tasks/doc_permission_syncing/tasks.pyelement_update_permissionsconnector_permission_sync_generator_task
索引后端迁移backend/onyx/background/celery/tasks/opensearch_migration/tasks.pymigrate_chunks_from_vespa_to_opensearch_task
迁移状态表backend/onyx/db/opensearch_migration.pyget_opensearch_migration_stateupdate_vespa_visit_progress_with_commit
定向重建backend/onyx/background/indexing/run_targeted_reindex.pyprocess_targets_for_cc_pair
写新连接器的规范backend/onyx/connectors/README.md

接着读哪一章: 想知道这些 chunk 被检索时怎么用,看 04 章;想知道 docfetching / docprocessing 这些 worker 怎么被调度、Redis 在里面协调什么,看 06 章