数据截至 (上游 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.py、factory.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:441 的 connector_document_extraction:它在 :762 把清洗后的批次写进 batch_storage,:788 用 app.send_task(OnyxCeleryTask.DOCPROCESSING_TASK, ...) 派发,然后 :815 保存 checkpoint。加工侧在 backend/onyx/background/celery/tasks/docprocessing/tasks.py:1853 调 run_indexing_pipeline。
这样拆的收益:连接器被源头限流时,不会占着嵌入 worker;嵌入慢时,也不会拖住抓取。worker 分工的全貌见 06 章。
3. 连接器契约:一份接口,53 个数据源
它要解决的小问题: 每个数据源的 API 都不一样,但下游管线只想拿到"一批 Document"。
思路: 不做一个万能大接口,而是拆成一组小的抽象基类,连接器按自己源头的能力挑着实现。instantiate_connector 出来的对象是什么类型,ConnectorRunner 就用 isinstance 走哪条路径。
3.1 四类主流:按"能不能增量、能不能续跑"分
| 契约 | 核心方法 | 语义 | 谁在用它 |
|---|---|---|---|
LoadConnector | load_from_state() | 全量:把源头当前完整状态吐一遍 | 首次索引 / 无法增量的源 |
PollConnector | poll_source(start, end) | 增量:只吐这个时间窗内变过的 | 有"最近修改"API 的源 |
SlimConnector | retrieve_all_slim_docs() | 只吐 ID,不吐正文 | 剪枝任务( 见 §7.1) |
CheckpointedConnector | load_from_checkpoint(start, end, ckpt) | 增量 + 断点续传 | 大源头的主力路径 |
定义分别在 backend/onyx/connectors/interfaces.py:120(LoadConnector)、:125(PollConnector)、:134(SlimConnector)、:266(CheckpointedConnector)。
两个"带权限同步"的变体是给企业版走 ACL 的:SlimConnectorWithPermSync.retrieve_all_slim_docs_perm_sync(:147)在拿 ID 的同时把每个文档的外部权限带回来;CheckpointedConnectorWithPermSync.load_from_checkpoint_with_perm_sync(:305)则在正常抓取时顺带把权限一起灌进 Document。
ConnectorRunner 会检查这一点:构造时如果 include_permissions=True 但连接器不是 CheckpointedConnector,直接抛 ValueError(backend/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 | 针对一批已知失败的文档重抓,不带 checkpoint | interfaces.py:316 |
EventConnector 是预留的——backend/onyx/connectors/README.md 里明说后台任务目前不用它。这是诚实的边界,不要以为 Onyx 已经支持 webhook 驱动索引。
Resolver 是 §7.4 定向重建的入口:它接收一组 ConnectorFailure,尽力把这些文档重新吐出来,由调用方负责把旧的失败记录换成新的。
3.3 基类里的公共动作
BaseConnector(interfaces.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_class(backend/onyx/connectors/factory.py:40)用 importlib.import_module 动态加载并缓存到 _connector_cache。好处很实在——不装 Salesforce SDK 的部署,也不会因为 import 失败而起不来。
实例化前还会校验"这个连接器支不支持你要的输入类型",_validate_connector_supports_input_type(factory.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_connector(factory.py:106)分两条路:
- 实现了
CredentialsConnector→ 注入一个OnyxDBCredentialsProvider,连接器自己按需读写凭据。 - 否则 → 直接解密一次
credential_json塞给load_credentials,如果连接器返回了新凭据就写回数据库。
OnyxDBCredentialsProvider(backend/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-35。LOCK_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:259 的 CheckpointOutput 别名)。三段式的含义:
- 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:54 的 CheckpointOutputWrapper 把上面那个 generator 转成统一的四元组流:
连接器 yield 的东西 包装后 yield 的四元组
─────────────────────────────────────────────────────
Document ──▶ (doc, None, None, None)
HierarchyNode ──▶ (None, node, None, None)
ConnectorFailure ──▶ (None, None, fail, None)
(最后的 return) ──▶ (None, None, None, checkpoint)
内部用一个 _inner_wrapper 做 self.next_checkpoint = yield from ... 把返回值截下来(:74-78),流结束后再单独 yield 一次。如果连接器忘了 return checkpoint,这里直接 RuntimeError(:90-93)——宁可炸也不要静默地丢掉存档点。
4.3 ConnectorRunner.run:四种连接器一个出口
ConnectorRunner.run(connector_runner.py:130)是"把四类契约折叠成一个接口"的地方。它做三件事:成批、补日志、统一出口形状(类注释 :99-104)。
对非 checkpoint 连接器,它伪造一个已完成的存档点:
finished_checkpoint = self.connector.build_dummy_checkpoint()
finished_checkpoint.has_more = False
connector_runner.py:229-230。之后 PollConnector 和 LoadConnector 各跑一遍自己的 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 拆成文档和节点两堆。