ARTICLE · 1066053
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。事故过程如下:
错误不只在于“旧向量没有删干净”。真正缺少的是可见性协议:当 v18 还未完整构建时,它不应对用户可见;v17 也不应被提前删除。
对任一
(tenant_id, doc_id),检索只能看到一个已完整构建且已发布的active_index_version;候选、失败和撤销版本不得参与召回。
二、事实源与版本模型
向量数据库是可重建的派生索引,不是事实源。PostgreSQL 决定文档是否存在、哪个源版本是期望版本、哪个索引版本已发布、租户与 ACL 是什么;对象存储提供不可变内容快照;向量库只做相似度检索。
不要把内容哈希称作“文档 ID”。至少要拆开如下身份:
doc_id | |
source_version | |
object_uri + etag | |
content_hash | |
index_version | |
embedding_modelchunker_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. 文件写入不可变对象路径,例如 tenant/acme/doc/doc-42/v18/source.pdf,绝不原地覆盖。2. 在一个 PostgreSQL 事务中写入 document_versions、更新文档的desired_source_version、插入 Outbox 记录。3. Outbox Relay 或 Debezium 将记录投递到 Kafka;Relay 崩溃会导致重复投递,但不会丢失事件。 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:minio, MINIO_ROOT_PASSWORD:minio123 }
ports: ["9000:9000"]
milvus:
image:milvusdb/milvus:v2.5.4
command: ["milvus", "run", "standalone"]
environment: { ETCD_ENDPOINTS:etcd:2379, MINIO_ADDRESS:minio:9000 }
depends_on: [etcd, minio]
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_id、doc_id、index_version、ordinal,让清理、审计和过滤有据可依。下面简化了解析和 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 端到端故障测试清单
以下用例是“可以照着实施”的最低验收线,而不仅是单元测试:
published_at 前退出 | ||
# 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 | |
desired != indexed | |
上线顺序应是:先完成稳定 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 配置参考