夜雨聆风学习资料网

ARTICLE · 1066053

5000 万向量、日均 300 次文档更新:RAG 增量一致性架构的实战复盘

5000 万向量、日均 300 次文档更新:RAG 增量一致性架构的实战复盘

5000 万向量、日均 300 次文档更新:RAG 增量一致性架构的实战复盘

关键词:RAG、增量索引、Transactional Outbox、Kafka、Milvus、版本发布、最终一致性

凌晨三点,客服机器人又报出了已经失效三个月的企业报价。报价 PDF 三天前已上传到对象存储,业务库也记录了新版本;但 embedding 服务超时后,索引任务只写入了一部分 chunk,旧 chunk 仍可检索。模型没有“记错”,它拿到的是一组过期且矛盾的证据。

这是 RAG Index Drift(索引漂移):事实源已改变,派生检索索引没有收敛。对于约 5000 万向量、日均约 300 次新增/修改/删除的知识库,全量重建一次需要 6~8 小时,不能作为日常同步方案。

本文以这次脱敏后的报价事故为案例,先解释设计,再给出一个可复制的最小实现。其目标不是让所有系统瞬时强一致,而是保证:失败、重复、乱序和删除发生后,用户只会检索到一个完整、可解释的文档版本。


一、贯穿案例:报价从 v17 更新到 v18 时发生了什么

文档 price-enterprise.pdf 的稳定业务主键为 doc-42。事故过程如下:

时间
事件
旧系统结果
周日 10:00
报价 v17 发布
10 个 v17 chunk 可检索
周一 09:05
上传 v18,价格变更
对象存储、业务库已是 v18
周一 09:06
embedding 超时
v18 仅写入 3 个 chunk
周一 09:07
任务先删除部分 v17
查询为空或混召新旧价格
周三 03:00
客服告警
索引漂移持续三天

错误不只在于“旧向量没有删干净”。真正缺少的是可见性协议:当 v18 还未完整构建时,它不应对用户可见;v17 也不应被提前删除。

对任一 (tenant_id, doc_id),检索只能看到一个已完整构建且已发布的 active_index_version;候选、失败和撤销版本不得参与召回。

二、事实源与版本模型

向量数据库是可重建的派生索引,不是事实源。PostgreSQL 决定文档是否存在、哪个源版本是期望版本、哪个索引版本已发布、租户与 ACL 是什么;对象存储提供不可变内容快照;向量库只做相似度检索。

不要把内容哈希称作“文档 ID”。至少要拆开如下身份:

字段
用途
doc_id
不随文件名或内容变化的业务主键
source_version
上传/删除事件的单调递增版本
object_uri + etag
Consumer 必须读取的不可变文件快照
content_hash
规范化 chunk 文本的去重与 embedding 缓存键
index_version
一次完整候选索引的发布版本
embedding_model
chunker_version
生成条件与兼容性边界

chunk_key = sha256(doc_id + index_version + ordinal + content_hash) 可作为 Milvus 主键。文本未变时,可以从 (content_hash, embedding_model) 缓存复用 embedding;但新发布版本仍应拥有独立的 chunk key,避免与旧版本互相覆盖。

v18 构建时,用户仍读完整 v17。只有 v18 的 chunk 数量、模型版本和 manifest 校验通过后,才把 active_index_version 从 17 原子切到 18;清理 v17 是后置、可重试的工作。

三、为什么必须使用 Outbox 和 CAS

Kafka producer transaction 不能同时覆盖对象存储、PostgreSQL 和 Kafka;它不是跨系统事务。因此上传流程应为:

  1. 1. 文件写入不可变对象路径,例如 tenant/acme/doc/doc-42/v18/source.pdf,绝不原地覆盖。
  2. 2. 在一个 PostgreSQL 事务中写入 document_versions、更新文档的 desired_source_version、插入 Outbox 记录。
  3. 3. Outbox Relay 或 Debezium 将记录投递到 Kafka;Relay 崩溃会导致重复投递,但不会丢失事件。
  4. 4. Consumer 允许重复处理,最终以条件更新(CAS)发布版本。

Kafka key 使用 tenant_id:doc_id,让同一文档事件大概率顺序读取。但不能把分区顺序当成正确性保证:重放、repair、消费者重平衡仍会产生旧事件;CAS 才是最后防线。

四、读路径、删除与对账也必须遵守可见性

只做到“先写新后删旧”仍可能让一次搜索召回两个版本。检索服务必须强制注入 tenant_id、ACL 与可见版本;不能信任调用方自行提供过滤条件。若向量库无法表达“每个 doc 不同 active version”,可用活动 collection/alias、分区切换,或召回后的严格可见性二次过滤。

删除也不能只调用 delete_by_doc_id:删除事务先写 deleted_at 和 DOCUMENT_DELETE_REQUESTED Outbox,检索服务立刻拒绝该文档,物理删除向量异步进行。权限收回采用同一策略,确保敏感内容即使仍占用向量库磁盘也不会被召回。

5000 万规模不能每小时将所有向量主键拉出做集合差。对账分三层:事件级(Outbox 延迟、Kafka lag、DLQ)、文档级(chunk count、manifest、版本状态)、分片级(按 tenant hash/时间窗口的摘要和分页深查)。Repair 事件必须复用正常 Consumer 和 CAS,不可绕过发布协议。


五、最小可运行实现:PostgreSQL + Kafka + Python + Milvus

以下实现专注于单租户/少量租户的最小闭环:上传后发布、重复事件幂等、候选版本完整校验、CAS 发布、读路径过滤、repair。生产环境再逐步添加认证、Kubernetes、对象存储、观测与分片。

5.1 目录与依赖

rag-indexer/
├── docker-compose.yml
├── .env.example
├── requirements.txt
├── sql/001_init.sql
├── app/
│   ├── config.py
│   ├── db.py
│   ├── schema.py
│   ├── outbox_relay.py
│   ├── consumer.py
│   ├── vector_store.py
│   └── api.py
└── tests/
    └── test_index_lifecycle.py
# requirements.txt
fastapi==0.115.6
uvicorn[standard]==0.34.0
sqlalchemy[asyncio]==2.0.36
asyncpg==0.30.0
aiokafka==0.12.0
pymilvus==2.5.4
pydantic-settings==2.7.0
pytest==8.3.4
pytest-asyncio==0.25.0

版本号应由项目锁文件固定;Milvus、Kafka 与 Python 客户端需要在 CI 中做兼容性验证。

5.2 本地基础设施

这里使用 Redpanda 兼容 Kafka API,以缩短本地启动路径;生产可替换为 Kafka 集群。Milvus standalone 仅用于开发与端到端测试,不代表 5000 万向量的生产拓扑。

# docker-compose.yml
services:
postgres:
image:postgres:16
environment:
POSTGRES_USER:rag
POSTGRES_PASSWORD:rag
POSTGRES_DB:rag
ports: ["5432:5432"]
healthcheck:
test: ["CMD-SHELL""pg_isready -U rag -d rag"]
interval:5s
timeout:3s
retries:20
kafka:
image:redpandadata/redpanda:v24.3.6
command:redpandastart--overprovisioned--smp1--memory1G--node-id0--check=false
ports: ["9092:9092"]
etcd:
image:quay.io/coreos/etcd:v3.5.16
command:etcd-advertise-client-urls=http://0.0.0.0:2379-listen-client-urls=http://0.0.0.0:2379
minio:
image:minio/minio:RELEASE.2024-10-29T16-01-48Z
command:minioserver/minio_data
environment: { MINIO_ROOT_USER:minioMINIO_ROOT_PASSWORD:minio123 }
ports: ["9000:9000"]
milvus:
image:milvusdb/milvus:v2.5.4
command: ["milvus""run""standalone"]
environment: { ETCD_ENDPOINTS:etcd:2379MINIO_ADDRESS:minio:9000 }
depends_on: [etcdminio]
ports: ["19530:19530"]
# .env.example
DATABASE_URL=postgresql+asyncpg://rag:rag@localhost:5432/rag
KAFKA_BOOTSTRAP_SERVERS=localhost:9092
KAFKA_TOPIC=rag.document.events
MILVUS_URI=http://localhost:19530
EMBEDDING_DIM=1536

启动后执行 docker compose up -d,再运行迁移。真实 embedding 服务的 API key 不应写进 Compose 或仓库;本地测试可替换为确定性的假 embedding。

5.3 数据库迁移

-- sql/001_init.sql
create type document_status as enum ('BUILDING''ACTIVE''DELETING''DELETED');
create type version_state as enum ('PENDING''BUILDING''READY''FAILED''DELETED');

create table documents (
  tenant_id text not null,
  doc_id uuid not null,
  desired_source_version bigintnot nulldefault0,
  indexed_source_version bigint,
  active_index_version bigint,
  status document_status not nulldefault'BUILDING',
  deleted_at timestamptz,
  updated_at timestamptz not nulldefault now(),
primary key (tenant_id, doc_id)
);

create table document_versions (
  tenant_id text not null,
  doc_id uuid not null,
  source_version bigintnot null,
  object_uri text not null,
  object_etag text not null,
  chunker_version text not null,
  embedding_model text not null,
  expected_chunk_count integer,
  manifest_hash text,
  state version_state not nulldefault'PENDING',
  error_message text,
  created_at timestamptz not nulldefault now(),
primary key (tenant_id, doc_id, source_version)
);

create table outbox (
  id bigserial primary key,
  aggregate_key text not null,
  event_type text not null,
  payload jsonb not null,
  created_at timestamptz not nulldefault now(),
  published_at timestamptz,
  attempts integernot nulldefault0
);
create index outbox_unpublished_idx on outbox (id) where published_at isnull;

数据库迁移工具可选 Alembic、Flyway 或 Liquibase;关键不在工具,而在于“更新业务状态和插入 Outbox”必须位于同一数据库事务。

5.4 上传 API:事务内写版本和 Outbox

该示例将对象上传抽象为 object_store.put_versioned。实现必须返回不可变 URI 和 ETag;禁止覆盖一个固定 key 后再让 Consumer 读取“最新对象”。

# app/api.py(关键路径)
@app.post("/documents/{doc_id}", status_code=202)
asyncdefupload(doc_id: UUID, file: UploadFile, tenant_id: str = Depends(tenant)):
    uri, etag = await object_store.put_versioned(tenant_id, doc_id, file)
asyncwith session.begin():
        version = await next_version(session, tenant_id, doc_id)
await session.execute(text("""
          insert into documents(tenant_id, doc_id, desired_source_version, status)
          values (:tenant, :doc, :version, 'BUILDING')
          on conflict (tenant_id, doc_id) do update
             set desired_source_version = excluded.desired_source_version,
                 status = 'BUILDING', deleted_at = null, updated_at = now()
        """
), {"tenant": tenant_id, "doc": doc_id, "version": version})
await session.execute(text("""
          insert into document_versions(tenant_id, doc_id, source_version, object_uri,
              object_etag, chunker_version, embedding_model)
          values (:tenant, :doc, :version, :uri, :etag, 'semantic-v2', 'embedding-v1')
        """
), {"tenant": tenant_id, "doc": doc_id, "version": version,
"uri": uri, "etag": etag})
await insert_outbox(session, tenant_id, doc_id, "DOCUMENT_UPSERT_REQUESTED",
            {"tenant_id": tenant_id, "doc_id"str(doc_id), "source_version": version,
"object_uri": uri, "etag": etag})
return {"doc_id"str(doc_id), "source_version": version, "status""BUILDING"}

5.5 Outbox Relay:允许重复,绝不静默丢失

Relay 认领少量未发布记录,使用聚合键作为 Kafka key。只有 Kafka broker 确认后才标记 published_at。崩溃发生在发送确认与更新之间时,消息会重发;这是预期行为。

# app/outbox_relay.py
asyncdefrelay_once(session, producer):
asyncwith session.begin():
        rows = (await session.execute(text("""
          select id, aggregate_key, event_type, payload
            from outbox where published_at is null
           order by id limit 100 for update skip locked
        """
))).mappings().all()
for row in rows:
await producer.send_and_wait(
                TOPIC, key=row["aggregate_key"].encode(),
                value=json.dumps({"type": row["event_type"], **row["payload"]}).encode(),
            )
await session.execute(text("""
              update outbox set published_at = now(), attempts = attempts + 1 where id = :id
            """
), {"id": row["id"]})

生产 Relay 应记录发送失败、最长未发布记录、重试次数与告警;若使用 Debezium,则以 WAL 捕获替代轮询,但 Consumer 幂等性仍然不可省略。

5.6 Milvus schema 与 Consumer:候选完整后才发布

向量记录至少存储 tenant_iddoc_idindex_versionordinal,让清理、审计和过滤有据可依。下面简化了解析和 embedding 的实现,但保留了决定一致性的控制流。

# app/vector_store.py(pymilvus 伪实现,字段与索引需与实际版本核对)
from pymilvus import DataType, MilvusClient

client = MilvusClient(uri=settings.milvus_uri)
COLLECTION = "rag_chunks"

defensure_collection():
if client.has_collection(COLLECTION):
return
    schema = client.create_schema(auto_id=False, enable_dynamic_field=False)
    schema.add_field("id", DataType.VARCHAR, is_primary=True, max_length=64)
    schema.add_field("vector", DataType.FLOAT_VECTOR, dim=settings.embedding_dim)
    schema.add_field("tenant_id", DataType.VARCHAR, max_length=64)
    schema.add_field("doc_id", DataType.VARCHAR, max_length=36)
    schema.add_field("index_version", DataType.INT64)
    schema.add_field("ordinal", DataType.INT32)
    index = client.prepare_index_params()
    index.add_index("vector", index_type="AUTOINDEX", metric_type="COSINE")
    client.create_collection(COLLECTION, schema=schema, index_params=index)

asyncdefcount(tenant_id, doc_id, version):
    expr = f'tenant_id == "{tenant_id}" and doc_id == "{doc_id}" and index_version == {version}'
returnlen(client.query(COLLECTION, filter=expr, output_fields=["id"], limit=16384))

生产代码不能直接用字符串拼接构造 filter,示例仅突出必需字段;应校验/转义输入,或使用 SDK 支持的参数化机制。对于 chunk 数超过单次查询限制的文档,需用分页或在 PostgreSQL 保存写入 manifest,不能只依赖上例的 limit

# app/consumer.py
asyncdefhandle_upsert(event, session, vectors):
# 只有 PENDING 或 FAILED 可以被认领;重复消息立即安全返回
    claimed = await session.execute(text("""
      update document_versions set state = 'BUILDING', error_message = null
       where tenant_id=:tenant and doc_id=:doc and source_version=:version
         and state in ('PENDING', 'FAILED') returning *
    """
), event)
    row = claimed.mappings().first()
ifnot row:
return

try:
        raw = await object_store.get_immutable(row["object_uri"], row["object_etag"])
        chunks = chunk_text(raw, row["chunker_version"])
        records = []
for ordinal, chunk inenumerate(chunks):
            content = normalize(chunk)
            records.append({"id": chunk_key(event["doc_id"], event["source_version"], ordinal, content),
"vector"await embed(content, row["embedding_model"]),
"tenant_id": event["tenant_id"], "doc_id": event["doc_id"],
"index_version": event["source_version"], "ordinal": ordinal})
await vectors.upsert(records)  # 相同主键覆盖,重复消费无重复向量
ifawait vectors.count(event["tenant_id"], event["doc_id"], event["source_version"]) != len(records):
raise RetryableError("candidate vector count mismatch")

# v19 已成为期望版本时,迟到 v18 不可覆盖 active 指针
        published = await session.execute(text("""
          update documents set active_index_version=:version, indexed_source_version=:version,
              status='ACTIVE', updated_at=now()
           where tenant_id=:tenant and doc_id=:doc
             and desired_source_version=:version and deleted_at is null
          returning doc_id
        """
), {"tenant": event["tenant_id"], "doc": event["doc_id"],
"version": event["source_version"]})
if published.first():
await session.execute(text("""
              update document_versions set state='READY', expected_chunk_count=:count
               where tenant_id=:tenant and doc_id=:doc and source_version=:version
            """
), {"tenant": event["tenant_id"], "doc": event["doc_id"],
"version": event["source_version"], "count"len(records)})
await insert_outbox(session, event["tenant_id"], event["doc_id"], "CLEANUP_OLD_VERSIONS",
                                {"keep": event["source_version"]})
except RetryableError:
await session.execute(text("""update document_versions set state='FAILED'
          where tenant_id=:tenant and doc_id=:doc and source_version=:version"""
), event)
raise# 不提交 Kafka offset;由消费框架重试

生产实现应将 embedding 改为批量调用,使用 (content_hash, model) 缓存,并通过重试主题或延迟队列控制退避。格式错误等不可恢复错误进入 DLQ,写明 doc_id/source_version/error_class;暂态错误不提交 offset。

5.7 检索和删除接口

读路径从 PostgreSQL 获取已发布、未删除的文档版本,再将其应用到向量过滤或召回后二次过滤。权限条件必须由服务端生成。

asyncdefretrieve(query: str, tenant_id: str, principal: Principal):
    visible = await metadata.visible_doc_versions(tenant_id, principal)
ifnot visible:
return []
    hits = await vectors.search(await embed_query(query), tenant_id=tenant_id, limit=50)
return [h for h in hits if visible.get(h.doc_id) == h.index_version]

asyncdefdelete_document(tenant_id: str, doc_id: UUID):
asyncwith session.begin():
        version = await next_version(session, tenant_id, doc_id)
await session.execute(text("""update documents set desired_source_version=:version,
          deleted_at=now(), status='DELETING' where tenant_id=:tenant and doc_id=:doc"""
),
          {"tenant": tenant_id, "doc": doc_id, "version": version})
await insert_outbox(session, tenant_id, doc_id, "DOCUMENT_DELETE_REQUESTED",
                            {"tenant_id": tenant_id, "doc_id"str(doc_id), "source_version": version})

visible_doc_versions 必须过滤 deleted_at is null 和权限。删除事件的 Consumer 删除/标记该文档的向量后,将状态推进为 DELETED;即便清理暂时失败,读路径已不可见。

5.8 端到端故障测试清单

以下用例是“可以照着实施”的最低验收线,而不仅是单元测试:

用例
注入方式
预期结果
重复投递 v18
同一 Kafka 消息投递两次
chunk 不重复,active 仍为 v18
v18 部分写入失败
第 3 批 embedding 抛超时
active 继续为 v17,v18 不可见
v18 迟到
先处理 v19,再投递 v18
CAS 拒绝 v18 发布,active 保持 v19
Relay 崩溃
Kafka 成功后、写 published_at 前退出
重发消息,最终仅一次可见发布
删除与向量库故障
删除后阻断 Milvus
查询立即无结果,清理可稍后完成
对账修复
人工删掉 active 中一个 chunk
manifest/count 检测异常并重建当前版本
# tests/test_index_lifecycle.py(核心断言)
asyncdeftest_partial_v18_never_replaces_v17(seed_v17, fail_embedding_on_call):
    fail_embedding_on_call(3)
await handle_upsert(v18_event, session, vectors)
assertawait active_version("acme", DOC_42) == 17
assertawait retrieve("企业报价""acme", sales_user) == await hits_for(17)

asyncdeftest_late_event_cannot_roll_back(seed_v19):
await handle_upsert(v18_event, session, vectors)
assertawait active_version("acme", DOC_42) == 19

六、生产验收与演进

不要用“日均更新数”直接决定分区数或 KEDA 副本,而要实测峰值文档大小、平均 chunk 数、embedding 吞吐、写入 P95 和可容忍的陈旧窗口。建议建立以下指标:

指标
建议目标
上传至 ACTIVE 的 P95
小于 10 分钟
Outbox 最老未发布记录
小于 1 分钟
Kafka lag 对应处理时长
小于 5 分钟
desired != indexed
 的文档数
稳态为 0
候选版本完整率
100% 才允许发布
删除/权限收回的检索拒绝延迟
小于 1 分钟

上线顺序应是:先完成稳定 ID、不可变对象与读路径过滤;再接 Outbox、幂等 Consumer、CAS 发布和 DLQ;随后加入 manifest 对账与自动 repair;最后根据 lag 和成本扩展 Kafka 分区、KEDA、embedding 缓存、Milvus 分片和跨区域部署。

embedding 模型、chunker、解析器、向量维度或权限 schema 升级时,创建新的 index_schema_version,构建候选索引并完成质量、容量和回滚验证后再切换可见版本。不要把新旧不兼容向量混在同一召回空间。

结语

RAG 增量一致性的本质是派生索引发布,而不是“新增几条向量、删除几条向量”。将事实源、Outbox、候选版本、CAS 发布、读路径过滤与分片对账连成同一个协议后,日均 300 次更新只处理真正变化的文本;失败和重试不会让文档消失,乱序也不能回退已发布版本。

最终,系统能够明确回答最重要的生产问题:此刻用户看到的是哪一版知识,为什么它有资格可见。

参考资料

  • • Spring for Apache Kafka:Transactions
  • • Apache Kafka Transaction Protocol
  • • Milvus DataCoord 配置参考

相关学习资料