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 记忆、知识图谱时,你几乎一定会撞到这几件事:
- 批处理 stale:跑全量 embedding 是几小时到几天的活,跑完就过期,再跑又浪费。
- 手工增量很贵:要识别新增 / 修改 / 删除,要跨 join 传播变化,要在向量库 / 关系库之间协调事务,逻辑量随数据量指数膨胀。
- 代码与数据同演化:换 chunk 策略、改 embed 模型、调 schema——所有这些都让「上次跑过的结果」失效,但你又不想整库重算。
- 多源异构: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 用 .env 设 COCOINDEX_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_embeddinghn_trending_topicsconversation_to_knowledgemulti_codebase_summarizationpatient_intake_extraction_bamlcsv_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. 坑与注意
@coco.fn(memo=True)的缓存粒度是 input + code:函数体里的非确定性调用(时间戳、随机数、外部 API 拉取)会让缓存键漂移或失效,要么放进参数显式传,要么别 memo。- 状态库与目标库是两件事:
COCOINDEX_DB是引擎自己用的 lineage 元数据库(默认 SQLite),与业务 Postgres 是分开的;生产别用默认 SQLite 跑长期。 - Target state 是「纯函数」:文档反复强调 target 没有副作用,所以同一份源 + 同一份代码必须产生同一份 target。如果你在 component 内部偷偷
requests.post(...)触发外部动作,那部分语义不会被引擎追踪,重跑时会重复触发——把它声明成另一类 component 才安全。 - Live update ≠ streaming:live update 仍由 CocoIndex 调度,靠源端 CDC 或定时触发,不替代 Kafka/Flink 的事件流。如果你要毫秒级延迟,先确认目标源有 CDC connector。
- chunk 策略改了会触发重算:这是设计行为,不是 bug;但大语料下第一次回填会重跑,建议先在样本上验证 chunker 再切生产。
- Postgres target 的向量索引后端:取决于 connector 配置,常见是 pgvector;如果你已经用了 Qdrant / Weaviate,需要选对应 connector,不要硬塞 pgvector。
- Python import 不要脑补:仓库 README 里给出的代码片段引用了
from cocoindex.connectors import localfs, postgres/from cocoindex.ops.text import RecursiveSplitter,都是公开 API,但子模块路径在不同版本可能调整,写完跑python -c "import cocoindex"验证再发布。 - 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 名称(
localfs、postgres、PatternFilePathMatcher、RecursiveSplitter、declare_vector_index)来自 README 与 quickstart,未逐条对照源码 API stability,使用前建议直接看docs/对应页。 - Live update 的具体 CDC connector 列表、官方提供的 sample 列表更新频率,本次仅以 README 当时截图为准。