用 CocoIndex 把 Markdown 文件夹变成可语义搜索的向量索引:分块、嵌入、pgvector 存储全流程
2026/9/15 22:02:28 网站建设 项目流程

用 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 中呈现,逻辑分四步:

  1. Walk:用localfs.walk_dir递归扫描本地目录(live=True开启变更监听能力);
  2. Chunk:用RecursiveSplitter把每个文件切成带重叠的片段——小而聚焦,重叠区保证跨边界的思想不会断成两截;
  3. Embed:用all-MiniLM-L6-v2嵌入每个片段——模型小巧快速、完全本地运行,无需 API Key;
  4. 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 embedder

DocEmbedding既是 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)) yield

coco.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包含textstartend三个字段,其中TextPosition携带byte_offsetchar_offsetlinecolumn四类位置信息,可用于精确回溯原文。

组装应用:挂载目标表与文件源

@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_targetmount_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_dirwalk_dir(path, live, recursive, path_matcher, rescan_interval)返回可异步迭代的DirWalker(python/cocoindex/connectors/localfs/_source.py)。live=True表示源支持实时监听,但真正进入 live 模式还需在 CLI 加-Lpath_matcherPatternFilePathMatcher(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 -d

2. 配置与安装

cp .env.example .env # set POSTGRES_URL (defaults to the local docker one) pip install -e .

.env.example 中可配置的环境变量:

变量说明默认值
POSTGRES_URLPostgres 连接串postgres://cocoindex:cocoindex@localhost/cocoindex
COCOINDEX_DBCocoIndex 本地元数据库路径./cocoindex.db
PYTORCH_ENABLE_MPS_FALLBACKMac 上 MPS 不支持的算子回退 CPU,其他平台无副作用1

依赖声明于 pyproject.toml:cocoindex[postgres,sentence_transformers]>=1.0.7asyncpg>=0.29.0pgvector>=0.4.1numpypython-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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询