用 CocoIndex 把 Markdown 文件夹变成可语义搜索的向量索引:分块、嵌入、pgvector 存储全流程
【免费下载链接】cocoindexIncremental engine for long horizon agents 🌟 Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex
导读
本指南以examples/text_embedding示例为核心,讲解如何用 CocoIndex 在纯异步 Python 中搭建一条「读取 Markdown → 递归分块 → 本地模型嵌入 → 写入 Postgres + pgvector」的端到端语义搜索管线:无需 API Key、无需手写增量逻辑,编辑一个文件只会重嵌一个文件。读完本文,你将掌握RecursiveSplitter分块、SentenceTransformerEmbedder嵌入、mount_table_target托管目标表与向量索引的完整用法,并能用一句自然语言查询召回「不含任何相同关键词」的正确段落。
背景:为什么需要向量索引
一堆 Markdown 文档里往往藏着答案,但传统的关键词检索(LIKE、全文索引)只能命中字面匹配——查询 "How does incremental processing work?" 无法召回一篇只讨论增量计算、措辞完全不同的文章段落。语义搜索的思路是:把文档切分成小块,每块用嵌入模型映射为稠密向量,查询时把问题嵌入成同一向量空间中的点,按余弦距离返回最近的块。这正是 RAG 与语义检索系统的共同地基。
CocoIndex 把这条流水线收敛为一行核心抽象:target_state = transformation(source_state)。数据变换用原生 Python 和自定义类型声明,而增量处理、变更追踪、托管目标(managed targets)等重活由底层 Rust 引擎承担,因此「改一个文件只重嵌一个文件,而不是重嵌整个文件夹」。
流水线总览:Walk → Chunk → Embed → Store
整条管线在 examples/text_embedding/main.py 中呈现,逻辑分四步:
- Walk:用
localfs.walk_dir递归扫描本地目录(live=True开启变更监听能力); - Chunk:用
RecursiveSplitter把每个文件切成带重叠的片段——小而聚焦,重叠区保证跨边界的思想不会断成两截; - Embed:用
all-MiniLM-L6-v2嵌入每个片段——模型小巧快速、完全本地运行,无需 API Key; - Store:每个片段在 Postgres 中写入一行,并为 embedding 列声明 pgvector 向量索引。
其中process_file对每个文件运行一次,memo=True让它具备增量能力:只要文件内容与函数代码均未变化,下次运行会整体跳过该文件。
示例自带三个样例文档(markdown_files/):1706.03762v7.md(Attention 论文)、1810.04805v2.md(BERT 论文)与rfc8259.md(JSON 规范),开箱即可体验。
核心代码逐段拆解
行类型:自己的 dataclass 就是表结构
@dataclass class DocEmbedding: id: int filename: str chunk_start: int chunk_end: int text: str embedding: Annotated[NDArray, EMBEDDER] # dimension inferred from the embedderDocEmbedding既是 Python 行对象,也是 Postgres 目标表的 schema 来源。embedding字段用Annotated[NDArray, EMBEDDER]标注——EMBEDDER是声明了detect_change=True的上下文键,向量维度由嵌入器自动推导(SentenceTransformerEmbedder实现VectorSchemaProvider,见 python/cocoindex/ops/sentence_transformers.py 中的__coco_vector_schema__,返回VectorSchema(dtype=float32, size=dim)),TableSchema.from_class据此生成 DDL。
生命周期:提供数据库连接池与嵌入器
@coco.lifespan async def coco_lifespan( builder: coco.EnvironmentBuilder, ) -> AsyncIterator[None]: async with asyncpg.create_pool(DATABASE_URL) as pool: builder.provide(PG_DB, pool) builder.provide(EMBEDDER, SentenceTransformerEmbedder(EMBED_MODEL)) yieldcoco.lifespan在应用启动/关闭时执行,向环境注入两个全局资源:asyncpg 连接池PG_DB,以及嵌入器EMBEDDER。注意EMBEDDER = coco.ContextKeySentenceTransformerEmbedder——detect_change=True意味着嵌入器的「身份」变化(如更换模型名)会被引擎识别为逻辑变更,从而自动对全量数据重嵌入,无需手动清缓存。
分块与嵌入:两个函数两层分工
@coco.fn async def process_chunk( chunk: Chunk, filename: pathlib.PurePath, id_gen: IdGenerator, table: postgres.TableTarget[DocEmbedding], ) -> None: table.declare_row( row=DocEmbedding( id=await id_gen.next_id(chunk.text), filename=str(filename), chunk_start=chunk.start.char_offset, chunk_end=chunk.end.char_offset, text=chunk.text, embedding=await coco.use_context(EMBEDDER).embed(chunk.text), ), ) @coco.fn(memo=True) async def process_file( file: FileLike, table: postgres.TableTarget[DocEmbedding], ) -> None: text = await file.read_text() chunks = _splitter.split( text, chunk_size=2000, chunk_overlap=500, language="markdown" ) id_gen = IdGenerator() await coco.map(process_chunk, chunks, file.file_path.path, id_gen, table)process_file负责读文件与分块,process_chunk负责逐块嵌入并声明行;id_gen.next_id(chunk.text)让id由块文本内容派生——这是增量更新的关键(见下文「变更如何被追踪」);chunk.start/end.char_offset记录了块在原文中的字符偏移,来自Chunk的位置信息(python/cocoindex/resources/chunk.py),便于回溯来源。
RecursiveSplitter 参数说明
RecursiveSplitter定义于 python/cocoindex/ops/text.py,支持语法感知的递归切分(markdown 走段落/句子边界,代码语言可走 tree-sitter)。核心参数:
| 参数 | 作用 | 示例取值 |
|---|---|---|
chunk_size | 目标块大小(字节,非字符数) | 2000 |
chunk_overlap | 相邻块之间重叠的字节数,防止跨边界语义断裂 | 500 |
min_chunk_size | 最小块大小,默认chunk_size / 2 | 不传则取1000 |
language | 语法感知语言名或扩展名(如"markdown"、"python"、".rs"),有 tree-sitter 支持时启用语法感知 | "markdown" |
返回的每个Chunk包含text、start、end三个字段,其中TextPosition携带byte_offset、char_offset、line、column四类位置信息,可用于精确回溯原文。
组装应用:挂载目标表与文件源
@coco.fn async def app_main(sourcedir: pathlib.Path) -> None: target_table = await postgres.mount_table_target( PG_DB, table_name=TABLE_NAME, table_schema=await postgres.TableSchema.from_class(DocEmbedding, primary_key=["id"]), pg_schema_name=PG_SCHEMA_NAME, ) target_table.declare_vector_index(column="embedding") files = localfs.walk_dir(sourcedir, recursive=True, path_matcher=PatternFilePathMatcher(included_patterns=["**/*.md"]), live=True) await coco.mount_each(process_file, files.items(), target_table) app = coco.App( coco.AppConfig(name="TextEmbeddingV1"), app_main, sourcedir=pathlib.Path("./markdown_files"), )mount_table_target:mount_table_target(db, table_name, table_schema, pg_schema_name)是「创建表目标 +coco.mount_target()挂载 + 包装」的组合糖(见 python/cocoindex/connectors/postgres/_target.py)。它全权托管表结构、幂等 upsert、以及源文件消失时的删除——你永远不需要写 diff 逻辑;declare_vector_index:声明 pgvector 索引,实际索引命名为{table_name}__vector__{name}。参数包括column(索引列)、metric("cosine"/"l2"/"ip",默认"cosine")、method("ivfflat"/"hnsw",默认"ivfflat"),以及 ivfflat 的lists和 hnsw 的m/ef_construction;walk_dir:walk_dir(path, live, recursive, path_matcher, rescan_interval)返回可异步迭代的DirWalker(python/cocoindex/connectors/localfs/_source.py)。live=True表示源支持实时监听,但真正进入 live 模式还需在 CLI 加-L;path_matcher用PatternFilePathMatcher(included_patterns=["**/*.md"])过滤,只收录 Markdown;mount_each:把process_file应用到文件流的每一项,并接到target_table上。
增量更新机制:为什么改一个文件只重嵌一个文件
@coco.fn(memo=True)为每个文件建立缓存:若文件内容与该函数的代码均未变化,下一轮运行直接跳过整个文件。而每行的id由块文本派生,因此重跑时:
- 内容未变的块 →
id不变 → 引擎幂等 upsert,结果不变; - 内容变化的块 →
id变化 → 更新对应行; - 源文件被删 → 相关块随之消失 → 引擎自动清理孤儿行。
此外EMBEDDER声明了detect_change=True,其 memo 键由(model_name_or_path, device, trust_remote_code)构成(见__coco_memo_key__),所以一旦更换嵌入模型,引擎会识别到缓存失效并对全量数据重新嵌入——这是「诚实的缓存失效」,无需手动清库。
查询:同一模型嵌入,余弦距离排序
示例查询代码同样位于 main.py:
async def query_once(pool, embedder, query, *, top_k=TOP_K) -> None: query_vec = await embedder.embed(query) async with pool.acquire() as conn: rows = await conn.fetch( f""" SELECT filename, text, embedding <=> $1 AS distance FROM "{PG_SCHEMA_NAME}"."{TABLE_NAME}" ORDER BY distance ASC LIMIT $2 """, query_vec, top_k, ) for r in rows: score = 1.0 - float(r["distance"]) print(f"[{score:.3f}] {r['filename']}") print(f" {r['text']}")- 查询时复用同一个
SentenceTransformerEmbedder(同一模型),保证索引与检索向量空间一致; <=>是 pgvector 的余弦距离运算符,score = 1 - distance换算为相似度(越接近 1 越相关),TOP_K = 5控制返回条数;- 交互模式支持反复输入查询(空行退出),也可一次传入参数直接查询。
运行指南
1. 启动 Postgres + pgvector
仓库提供了现成的 compose 配置 dev/postgres.yaml(镜像pgvector/pgvector:pg17,账号/密码/库名均为cocoindex,端口5432):
docker compose -f ../../dev/postgres.yaml up -d2. 配置与安装
cp .env.example .env # set POSTGRES_URL (defaults to the local docker one) pip install -e ..env.example 中可配置的环境变量:
| 变量 | 说明 | 默认值 |
|---|---|---|
POSTGRES_URL | Postgres 连接串 | postgres://cocoindex:cocoindex@localhost/cocoindex |
COCOINDEX_DB | CocoIndex 本地元数据库路径 | ./cocoindex.db |
PYTORCH_ENABLE_MPS_FALLBACK | Mac 上 MPS 不支持的算子回退 CPU,其他平台无副作用 | 1 |
依赖声明于 pyproject.toml:cocoindex[postgres,sentence_transformers]>=1.0.7、asyncpg>=0.29.0、pgvector>=0.4.1、numpy、python-dotenv>=1.0.1,要求 Python ≥ 3.11。
3. 构建索引
cocoindex update main # catch-up: scan, sync, exit cocoindex update -L main # live: keep watching for file changes第一条为一次性追平(扫描、同步、退出);第二条为 live 模式,持续监听文件变化并自动增量同步。
4. 语义搜索
python main.py "what is self-attention?"查询被同一模型嵌入后按余弦距离召回最相似的块并排序——即使它们与查询不含任何相同单词。这就是向量索引的意义。
从源码看底层保障
- 嵌入批处理与 OOM 保护:
SentenceTransformerEmbedder._embed以@coco.fn.as_async(batching=True, runner=coco.GPU, max_batch_size=64)声明,并发单文本调用会被引擎自动合批,GPU 显存不足时抛出coco.RetryWithSmallerBatch()并清空加速器缓存后减半重试(见 python/cocoindex/ops/sentence_transformers.py); - 目标表声明语义:
TableTarget.declare_row提取主键列值并声明目标状态(upsert 语义);declare_vector_index通过 attachment 机制管理索引生命周期——表或索引变更时由引擎统一协调建/删/重建(python/cocoindex/connectors/postgres/_target.py); - 文件监听的自愈:
walk_dir(live=True)默认每小时全量重扫一次并重建 OS 级 watcher(rescan_interval可调、可禁用),防御平台级 watcher 失效(如 macOS FSEvents 静默停止)。
小结
examples/text_embedding是 CocoIndex 语义检索场景的「最小完整闭环」:一个main.py覆盖 walk → chunk → embed → store → search 全部环节,增量处理、变更追踪与托管目标由 Rust 引擎透明承担。要扩展为真正的 RAG 应用,只需替换嵌入模型(HuggingFace 上任意 sentence-transformers 模型均可)、调整分块参数,或把目标换成其他受支持的存储即可。
【免费下载链接】cocoindexIncremental engine for long horizon agents 🌟 Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考