cocoindex-io/cocoindex · 上手攻略

  • 仓库:cocoindex-io/cocoindex
  • 链接:https://github.com/cocoindex-io/cocoindex · 文档 https://cocoindex.io/docs · Discord https://discord.com/invite/zpA9S2DR7s
  • 分类:AI / 数据索引 / ETL 引擎
  • 作者:spark
  • 更新:2026-09-06

1. 是什么

CocoIndex 是一个面向 AI 工作负载的增量数据索引引擎:用声明式 Python 描述「源数据经过若干变换后应该长成什么样」,引擎自动算出 Δ(变化的部分),把目标状态(Postgres、向量库、本地文件、对象存储等)始终保持最新——而无需你手写 diff、批量调度、状态恢复或回填脚本。

核心叙事:把 React 的「UI = f(state)」心智模型套到数据流水线上——你声明 TargetState = Transform(SourceState),引擎负责把任何 source 变化映射成正确的 target 增量。文档里把它类比为电子表格的公式、React 的渲染、或 SQL 的物化视图。

底层是 Rust 实现(核心调度、并行分块、零拷贝变换、故障隔离),上层 Python SDK 友好,官方口号「Python not a DAG」,意指不需要画 DAG、不需要写 Airflow。

2. 解决什么问题

构建 RAG、agent 记忆、知识图谱时,你几乎一定会撞到这几件事:

  1. 批处理 stale:跑全量 embedding 是几小时到几天的活,跑完就过期,再跑又浪费。
  2. 手工增量很贵:要识别新增 / 修改 / 删除,要跨 join 传播变化,要在向量库 / 关系库之间协调事务,逻辑量随数据量指数膨胀。
  3. 代码与数据同演化:换 chunk 策略、改 embed 模型、调 schema——所有这些都让「上次跑过的结果」失效,但你又不想整库重算。
  4. 多源异构:Slack/Notion/GitHub/邮件/PDF/视频会议记录要统一成 RAG-friendly 的 chunk+embedding,传统 ETL 不擅长文本语义切分。

CocoIndex 把这四件事一锅端:用 @coco.fn(memo=True) 做函数级 memoization(输入和代码都没变就跳过),用 declare_* 系列声明 target state,用 mount_each 把源条目映射到 processing component,由引擎在后台做 lineage 追踪、差量推送、目标回滚。

3. 快速安装

环境:Python 3.10+、pip 或 uv。文档当前版本 v1.0.14(2026-07-04 校对)

# 基础包
pip install -U cocoindex

# 常用组合:本地文件 + Postgres + 文本切分
pip install -U cocoindex psycopg2-binary

CLI 入口 cocoindex 在安装后即可用,所有 pipeline 都通过 cocoindex update <file>.py 触发。

⚠️ 默认需要一个轻量状态库(CocoIndex 自己维护的 SQLite/Postgres 元数据库)记录 lineage 和增量状态。本地试用可走默认 SQLite;接 Postgres target 时建议把状态库也指向 Postgres。Quickstart 用 .envCOCOINDEX_DB=./cocoindex.db

4. 核心用法

4.1 心智模型三件套

  • App:顶层可执行单元,把主函数绑到具体输入输出。
  • Processing Component:一个最小可独立处理单元(一个文件、一行记录、一段对话),内部声明自己的 target state。
  • Function(@coco.fn:变换函数;memo=True 开启按 hash(input)+hash(code) 缓存。

4.2 最小例子:把本地 PDF 转 Markdown(来自官方 Quickstart)

# main.py
import pathlib
import cocoindex as coco
from cocoindex.connectors import localfs
from cocoindex.resources.file import PatternFilePathMatcher
from docling.document_converter import DocumentConverter

_converter = DocumentConverter()

@coco.fn(memo=True)
def process_file(file: localfs.File, outdir: pathlib.Path) -> None:
    md = _converter.convert(file.file_path.resolve()).document.export_to_markdown()
    localfs.declare_file(
        outdir / (file.file_path.path.stem + ".md"),
        md,
        create_parent_dirs=True,
    )

@coco.fn
async def app_main(sourcedir: pathlib.Path, outdir: pathlib.Path) -> None:
    files = localfs.walk_dir(
        sourcedir,
        recursive=True,
        path_matcher=PatternFilePathMatcher(included_patterns=["**/*.pdf"]),
    )
    await coco.mount_each(process_file, files.items(), outdir)

app = coco.App("PdfToMarkdown", app_main,
               sourcedir=pathlib.Path("./pdf_files"),
               outdir=pathlib.Path("./out"))

运行:

pip install -U cocoindex docling
mkdir cocoindex-quickstart && cd cocoindex-quickstart
mkdir pdf_files   # 放 PDF
echo "COCOINDEX_DB=./cocoindex.db" > .env
# 上面 main.py 落盘后
cocoindex update main.py

增删 PDF 之后再次 cocoindex update,只重算变化的源条目;删除文件会让对应输出被自动回收。

4.3 经典 RAG:本地文档 → chunk → embedding → Postgres + 向量索引

README 主例(核心骨架):

import cocoindex as coco
from cocoindex.connectors import localfs, postgres
from cocoindex.ops.text import RecursiveSplitter

@coco.fn(memo=True)   # ← 按 input + code 双哈希缓存
async def index_file(file, table):
    for chunk in RecursiveSplitter().split(await file.read_text()):
        table.declare_row(text=chunk.text, embedding=embed(chunk.text))

@coco.fn
async def main(src):
    table = await postgres.mount_table_target(PG, table_name="docs")
    table.declare_vector_index(column="embedding")
    await coco.mount_each(index_file, localfs.walk_dir(src).items(), table)

coco.App(coco.AppConfig(name="docs"), main, src="./docs").update_blocking()

要点:

  • table.declare_row(...) 声明一行 target state,自动维护主键去重;
  • table.declare_vector_index(column=...) 让 Postgres target 自动创建 pgvector / 外部向量索引(具体后端依赖 connector);
  • mount_each(index_file, items, table) 把每个源条目映射成一个 processing component,组件内部可多次 declare。

4.4 实时增量(Live update)

对支持 CDC 的源(Postgres、Kafka、S3 event 等),用 coco.App(...).update() 之外的形式挂 listener 即可持续同步;非 CDC 源用 cocoindex update <file> 反复触发,引擎自带 dedupe + 版本控制。详见 docs programming_guide/app

4.5 仓库自带的可复制样例

examples/ 下 20+ 个官方 starter,按周更新,README 已列:

  • code_embedding / pdf_embedding
  • hn_trending_topics
  • conversation_to_knowledge
  • multi_codebase_summarization
  • patient_intake_extraction_baml
  • csv_to_kafka

每个样例都是「clone → 填连接串 → 跑」的形态,对照改 connector 参数最快上手。

4.6 与 AI 编程 agent 协作

仓库自带 skills/cocoindex skill(一个面向 Claude/Cursor 的 SKILL.md),让 agent 一次性读到编程模型 + API + 模式,再写 v1 代码时减少试错。官方推荐用法:在 agent 项目里挂这个 skill,然后让 agent 直接写 App。

5. 典型适用场景

  • RAG 索引:本地文档/Notion/Slack/邮件/PDF → chunk + embedding → Postgres+pgvector 或 Qdrant,源变化秒级同步。
  • Agent 长期记忆:把会话、ticket、wiki 转成结构化 memory,定期 incremental 入库(对话 → 知识样例即此类)。
  • 代码库语义检索code_embedding 样例,多仓汇总 + 多粒度切分。
  • 结构化抽取patient_intake_extraction_baml 用 BAML 做 schema-aware 抽取,结果落 Postgres。
  • 流式数据入湖csv_to_kafka 把 CSV 转 Kafka 事件,continuous 模式。
  • 视觉文档检索:配合 ColPali 做多向量 patch 检索(官方博客示例)。

不太适合:纯批处理一次性报表、必须用 SQL 表达的全部业务逻辑(直接用 dbt / sqlmesh 更顺手)、强 OLAP 复杂 join(这部分用 warehouse 联邦查询)。

6. 坑与注意

  1. @coco.fn(memo=True) 的缓存粒度是 input + code:函数体里的非确定性调用(时间戳、随机数、外部 API 拉取)会让缓存键漂移或失效,要么放进参数显式传,要么别 memo。
  2. 状态库与目标库是两件事COCOINDEX_DB 是引擎自己用的 lineage 元数据库(默认 SQLite),与业务 Postgres 是分开的;生产别用默认 SQLite 跑长期。
  3. Target state 是「纯函数」:文档反复强调 target 没有副作用,所以同一份源 + 同一份代码必须产生同一份 target。如果你在 component 内部偷偷 requests.post(...) 触发外部动作,那部分语义不会被引擎追踪,重跑时会重复触发——把它声明成另一类 component 才安全。
  4. Live update ≠ streaming:live update 仍由 CocoIndex 调度,靠源端 CDC 或定时触发,不替代 Kafka/Flink 的事件流。如果你要毫秒级延迟,先确认目标源有 CDC connector。
  5. chunk 策略改了会触发重算:这是设计行为,不是 bug;但大语料下第一次回填会重跑,建议先在样本上验证 chunker 再切生产。
  6. Postgres target 的向量索引后端:取决于 connector 配置,常见是 pgvector;如果你已经用了 Qdrant / Weaviate,需要选对应 connector,不要硬塞 pgvector。
  7. Python import 不要脑补:仓库 README 里给出的代码片段引用了 from cocoindex.connectors import localfs, postgres / from cocoindex.ops.text import RecursiveSplitter,都是公开 API,但子模块路径在不同版本可能调整,写完跑 python -c "import cocoindex" 验证再发布。
  8. CLI 命令是 cocoindex update,不是 coco run / coco exec——这是新用户最常见的踩坑。

7. 与同类对比

工具 核心定位 与 CocoIndex 的差异
Airflow / Prefect / Dagster 通用 workflow 编排 它们调度 DAG,CocoIndex 不画 DAG、不调度,靠 lineage 算增量;适合放进 DAGster 内部当 transform 步骤
dbt SQL 仓内转换 dbt 强在 SQL 表达,CocoIndex 强在跨异构源 + 任意 Python transform + 向量/文件 target
Unstructured / LlamaIndex ingestion 文档 → chunk 工具 它们解决「怎么切」,CocoIndex 解决「切完怎么持续同步 + 多源 join」
LangChain / LlamaIndex RAG LLM 编排框架 它们负责「检索后怎么喂给 LLM」,CocoIndex 负责「被检索的那份索引怎么保鲜」
Spark / Flink 大数据流批 Spark/Flink 是通用执行引擎,CocoIndex 是面向 AI 数据语义的专用引擎,体量小、上手快,但跨 PB 级 ETL 不替代 Spark
Mage / Prefect pipelines notebook 化 pipeline 同上调度思路,但缺 lineage 增量 + 多 target 原子写入
River / Kafka Streams 流式增量 真流式;CocoIndex 增量是 batch-friendly,对低延迟场景要靠 CDC connector 桥接

一句话:CocoIndex 不是要替代 ETL/编排平台,而是把「RAG/agent 索引保鲜」这一类高频痛点单独做成产品。

8. 一句话推荐结论

如果你正在做 RAG / agent 长期记忆 / 知识图谱,并且源数据每天都在变——直接用 CocoIndex,省下的是「写增量逻辑 + 写状态恢复 + 写 target 一致性」整套工程债;如果你只是做一次性离线批处理或者强 OLAP,沿用 dbt/Spark 更稳。


不确定 / 已显式标注

  • 文档版本 v1.0.14(页面标注 2026-07-04 校对),但未独立核对 PyPI 当前 latest 版本号;请生产切版本前 pip index versions cocoindex 复核。
  • 部分 connector 名称(localfspostgresPatternFilePathMatcherRecursiveSplitterdeclare_vector_index)来自 README 与 quickstart,未逐条对照源码 API stability,使用前建议直接看 docs/ 对应页。
  • Live update 的具体 CDC connector 列表、官方提供的 sample 列表更新频率,本次仅以 README 当时截图为准。