Streaming-Online-RAG

预印本勘误后全量重建索引,线上仍在回答旧方法。根因是索引生命周期按批而不是按变更事件。本方法做增量、幂等与可见性控制。一致性与回滚比吞吐更难,必须报 Time-to-Searchable。

本文属于 RAG 工程框架中的「5 在线运营与成本治理」环节,聚焦「Streaming Online RAG」方法。

定位

维度 内容
角色 索引的增量物化
输入 → 输出 变更事件 → 可检索的新版本 chunk
默认组合 Kafka + Debezium;LlamaIndex ingestion 增量
何时不用 静态库按月更新

核心机制

$$
\mathrm{TTS}=t_{\mathrm{searchable}}-t_{\mathrm{event}}
$$

写入按 doc_id+version 幂等。

实现路径与心智:变更事件进队列,按版本幂等 upsert 嵌入与可见性,失败可回滚。底层心智:索引是变更日志的物化视图,不是每月打一次的快照。

优缺点

  • 优点:新鲜度高。
  • 缺点:一致性与回滚复杂。

契约与走通样例

输入

1
{"event": "update", "doc_id": "P-GAPDH-01", "version": 3, "patch": "48 h → 24 h (erratum)"}

中间量

v3 upsert 后旧 48 h 块不可见。TTS=45 s。查询已命中 24 h。

输出

1
{"visible_version": 3, "tts_s": 45}

社区实现

Kafka + Debezium。风险:可见性与嵌入提交非原子,会出现「搜到空向量」。

工程落地

最小可运行示例

复制为 .py 后直接运行(仅标准库)。生产 upsert 对应向量库;必须嵌入完成后再切可见版本。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
"""幂等写入新版本,publish 之前对外不可见。"""
from __future__ import annotations

from dataclasses import dataclass, field


@dataclass
class StreamIndex:
"""版本化索引。生产对应 Qdrant payload.version + alias 切换。"""

visible: dict[str, str] = field(default_factory=dict)
staging: dict[tuple[str, str], list[str]] = field(default_factory=dict)

def upsert(self, doc_id: str, version: str, chunks: list[str]) -> None:
self.staging[(doc_id, version)] = chunks

def publish(self, doc_id: str, version: str) -> None:
if (doc_id, version) not in self.staging:
raise KeyError("embed unfinished")
self.visible[doc_id] = version


def ingest_stream(doc_id: str, version: str, text: str, index: StreamIndex) -> str:
"""输入事件;输出对外可见 version。"""
chunks = [p.strip() for p in text.split(".") if p.strip()]
index.upsert(doc_id, version, chunks)
index.publish(doc_id, version)
return version


if __name__ == "__main__":
index = StreamIndex()
ingest_stream("P-GAPDH-01", "v2", "Treat with 20 nM siRNA. Extract RNA at 48 h.", index)
ingest_stream("P-GAPDH-01", "v2", "Treat with 20 nM siRNA. Extract RNA at 48 h.", index) # 幂等
print(index.visible, index.staging[("P-GAPDH-01", "v2")])

参数

参数 起点 影响
TTS SLA 分钟级 过紧则要跳过重解析
幂等键 doc_id+version 缺则重复嵌入

失效—信号—螺丝

  • 搜到旧版:螺丝:发布前先下线旧 version。
  • 重复事件:螺丝:幂等键。
  • 半发布:螺丝:嵌入完成再切可见性。

规模(100 篇生物学 PDF)

全量重放约 1–5 h;端到端 TTS 常 秒–分钟级。嵌入仍需 GPU 8–16 GB

-------------本文结束感谢您的阅读-------------