news 2026/9/15 21:22:11

CocoIndex 管线模式实战指南:从文件转换到 LLM 抽取的六种增量处理范式

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
CocoIndex 管线模式实战指南:从文件转换到 LLM 抽取的六种增量处理范式

CocoIndex 管线模式实战指南:从文件转换到 LLM 抽取的六种增量处理范式

【免费下载链接】cocoindexIncremental engine for long horizon agents 🌟 Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex

CocoIndex 是一款 Python 原生的增量数据处理引擎,其核心设计理念是"声明目标状态,而非编写更新逻辑"——你只需描述"输出应该长什么样",引擎会自动完成变更检测、增量同步与外部系统的增删改。本文基于skills/cocoindex/references/patterns.md中的通用模式文档,系统讲解文件转换、向量嵌入、数据库 ETL、LLM 结构化抽取、Kafka 流式处理与共享资源上下文管理这六种最常见的管线范式,并结合仓库源码与可运行示例给出深度实现佐证。读完本文,你将掌握如何用@coco.fnmount_each()ContextKey等核心 API 搭建可增量更新、可自动清理、可实时监听的生产级数据管线。

核心范式:TargetState = Transform(SourceState)

所有 CocoIndex 管线都遵循同一条声明式三步流程:

  1. 读取源状态(Read source state):从文件系统、数据库、Kafka 等源读取当前数据;
  2. 变换(Transform):用@coco.fn装饰的纯处理函数完成格式转换、分块、嵌入、LLM 抽取等逻辑;
  3. 声明目标状态(Declare target state):调用declare_row()declare_file()declare_target_state()等 API 声明"目标应该存在什么"。

增量同步由 CocoIndex 引擎自动完成:它记录每次运行的目标状态,在下一次运行时通过 reconcile 机制对比"期望状态"与"当前状态",自动生成创建、更新或删除动作。以本地文件系统目标为例,localfs 目标实现中的_reconcile_entry()会对每个条目计算内容指纹(fingerprint_bytes),只有指纹与上次记录不一致时才触发写入;删除动作则通过NonExistenceType表达——源文件被删除后,目标文件会被自动清理,这正是"自动清理(Auto-cleanup)"的底层原理。

该范式与 React 的"声明式 UI"异曲同工:你声明"输入是什么,输出就应是什么",框架负责把状态收敛到声明的结果,TargetState = Transform(SourceState)正是 CocoIndex Skill 文档 中反复强调的第一原则。

Pattern 1:文件转换管线(File Transformation Pipeline)

适用场景:将文件从一种格式转换为另一种格式,例如 Markdown → HTML、PDF → Markdown、文本清洗等。

完整代码示例

import pathlib import cocoindex as coco from cocoindex.connectors import localfs from cocoindex.resources.file import FileLike, PatternFilePathMatcher from markdown_it import MarkdownIt _markdown_it = MarkdownIt("gfm-like") @coco.fn(memo=True) async def process_file(file: FileLike, outdir: pathlib.Path) -> None: html = _markdown_it.render(await file.read_text()) outname = "__".join(file.file_path.path.parts) + ".html" localfs.declare_file(outdir / outname, html, create_parent_dirs=True) @coco.fn async def app_main(sourcedir: pathlib.Path, outdir: pathlib.Path) -> None: files = localfs.walk_dir( sourcedir, path_matcher=PatternFilePathMatcher(included_patterns=["**/*.md"]), live=True, # Enable live file watching ) await coco.mount_each(process_file, files.items(), outdir) app = coco.App( coco.AppConfig(name="FilesTransform"), app_main, sourcedir=pathlib.Path("./data"), outdir=pathlib.Path("./output_html"), )

关键点剖析

  • memo=True—— 跳过未变化文件的重复处理@coco.fn(memo=True)开启函数级记忆化。从 function.py 的实现看,记忆化的指纹由两部分构成:调用指纹(输入参数经过fingerprint_call()规范化)与逻辑指纹(函数源码 AST 规范化后的哈希,见_compute_logic_fingerprint())。只有两者都未变化时才直接复用上次结果,从而跳过对未变化文件的读取、转换与写入。特别值得一提的是,逻辑指纹基于 AST 而非原始文本,因此注释、空行、格式调整不会引起误判失效;显式传入version=参数则会强制使指纹变化、触发重跑。

  • mount_each()—— 每个文件一个组件mount_each(fn, items, *args)items中的每一项挂载一个相互独立的处理组件,items(key, value)键值对,key 自动成为组件子路径(component subpath)。在 api.py 的实现中,mount_each()对每个 key 调用item_path = child_path.concat(key)构造唯一路径。文件转换中files.items()返回(relative_path, File)对(见 localfs 源实现 的items()),因此每个文件的组件路径就是其相对路径——这是稳定的、可记忆化的身份标识。

  • live=True—— 实时文件监听walk_dir(..., live=True)时,items()返回一个LiveMapView。底层使用watchdogObserver注册操作系统级监听(inotify / FSEvents),先做一次全量扫描(subscriber.update_all()),随后持续接收created / modified / deleted / moved事件并转换为增量更新。moved事件会被拆解为"删除旧路径 + 创建新路径"。作为防御机制,默认每小时做一次周期性全量重扫(rescan_interval,默认 1 小时),用于兜底平台级监听失效(如 macOS FSEvents 静默停止),可通过rescan_interval=None关闭。

  • 自动清理:由于每个输出文件都是一个声明式目标状态,当源文件被删除、对应组件不再声明目标时,引擎会依据 reconcile 结果自动删除对应的输出文件,无需手写清理逻辑。

仓库中的 files_transform 示例 与此模式完全对应,是开箱即用的参考起点。

Pattern 2:向量嵌入管线(Vector Embedding Pipeline)

适用场景:对文档做分块(chunking)并生成向量嵌入,写入向量数据库用于语义搜索。

完整代码示例

import pathlib from dataclasses import dataclass from typing import AsyncIterator, Annotated import asyncpg from numpy.typing import NDArray import cocoindex as coco from cocoindex.connectors import localfs, postgres from cocoindex.ops.text import RecursiveSplitter from cocoindex.ops.sentence_transformers import SentenceTransformerEmbedder from cocoindex.resources.chunk import Chunk from cocoindex.resources.file import FileLike, PatternFilePathMatcher from cocoindex.resources.id import IdGenerator DATABASE_URL = "postgres://cocoindex:cocoindex@localhost/cocoindex" PG_DB = coco.ContextKeyasyncpg.Pool EMBEDDER = coco.ContextKeySentenceTransformerEmbedder _splitter = RecursiveSplitter() @dataclass class DocEmbedding: id: int filename: str chunk_start: int chunk_end: int text: str embedding: Annotated[NDArray, EMBEDDER] @coco.lifespan async def coco_lifespan(builder: coco.EnvironmentBuilder) -> AsyncIterator[None]: async with await asyncpg.create_pool(DATABASE_URL) as pool: builder.provide(PG_DB, pool) builder.provide(EMBEDDER, SentenceTransformerEmbedder("all-MiniLM-L6-v2")) yield @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) @coco.fn async def app_main(sourcedir: pathlib.Path) -> None: target_table = await postgres.mount_table_target( PG_DB, table_name="doc_embeddings", table_schema=await postgres.TableSchema.from_class(DocEmbedding, primary_key=["id"]), ) target_table.declare_vector_index(column="embedding") files = localfs.walk_dir(sourcedir, recursive=True, path_matcher=PatternFilePathMatcher(included_patterns=["**/*.md"])) await coco.mount_each(process_file, files.items(), target_table) app = coco.App(coco.AppConfig(name="TextEmbedding"), app_main, sourcedir=pathlib.Path("./markdown_files"))

关键点剖析

  • mount_table_target(PG_DB, ...)ContextKey作为第一个参数:表目标通过db这个ContextKey定位连接池。在 postgres 目标实现 中,table_target()_TableKey(db_key=db.key, pg_schema_name, table_name)构造目标状态的全局唯一键——ContextKey的字符串 key 就是目标身份的稳定标识,这也是后文"上下文管理"模式强调key不能随意改名的原因。mount_table_target()table_target()+coco.mount_target()的语法糖。

  • Annotated[NDArray, EMBEDDER]—— 向量维度自动推断DocEmbedding.embeddingAnnotated标注了ContextKey,引擎可从SentenceTransformerEmbedder("all-MiniLM-L6-v2")推断出向量维度,并据此在 PostgreSQL 中创建对应的pgvector列。

  • map()—— 组件内并发执行,不创建子组件:与mount_each()不同,coco.map(fn, items, ...)是纯 asyncio 并发执行(TaskGroup),不产生组件子路径,因此不参与记忆化与目标状态的独立管理。从 api.py 的实现看,map()会等待所有任务完成,且按输入顺序抛出首个失败(DeadlineExceededError会被单独识别并重抛)。在向量管线的用法中,process_file是整个组件的记忆化边界(memo=True),其内部的分块与嵌入通过map()并发完成——文件不变则整组分块全部跳过。

  • IdGenerator—— 跨增量更新保持稳定的唯一 ID:分块场景中两个块可能内容相同,generate_id(dep)(同 dep 同 ID)会冲突,因此使用IdGenerator。其实现(见 id.py)为每个dep维护一个自增序号(ordinal),next_id()实际调用一个以(deps_fp, dep_fp, ordinal)为键的@coco.fn(memo=True)内部函数,从而保证"每次调用返回不同 ID、且跨运行序列稳定"——即使源文件内容未变,之前分配的 ID 也不会漂移,这是向量库主键稳定的关键。

  • RecursiveSplitter分块参数split(text, chunk_size=2000, chunk_overlap=500, language="markdown")指定块大小、重叠量与语言。该实现位于 text.py,底层是 Rust 实现的递归分隔器;language还支持通过detect_code_language(filename=...)按文件扩展名自动检测编程语言,CustomLanguageConfig可自定义分隔规则。

  • declare_vector_index(column="embedding"):声明 pgvector 索引作为表的附件目标状态。签名支持metriccosine/l2/ip,默认cosine)、methodivfflat/hnsw,默认ivfflat)、listsmef_construction等参数,实际创建的索引命名为{table_name}__vector__{name}

  • memo=Trueprocess_file:文件未变化时整块跳过,避免重复读取、分块与嵌入——这是控制成本的命脉。

仓库中的 text_embedding 示例 与本模式完全一致,另见 text_embedding_lancedb 等以不同向量库为目标的变体。

Pattern 3:数据库源 → 变换 → 数据库目标

适用场景:将一张数据库表变换后同步到另一张表,典型 ETL。

完整代码示例

from dataclasses import dataclass from typing import AsyncIterator import asyncpg import cocoindex as coco from cocoindex.connectors import postgres SOURCE_DB_URL = "postgres://localhost/source_db" TARGET_DB_URL = "postgres://localhost/target_db" SOURCE_DB = coco.ContextKeyasyncpg.Pool TARGET_DB = coco.ContextKeyasyncpg.Pool @dataclass class SourceRecord: id: int name: str value: float @dataclass class TargetRecord: id: int name: str value: float processed: bool @coco.lifespan async def coco_lifespan(builder: coco.EnvironmentBuilder) -> AsyncIterator[None]: async with ( await asyncpg.create_pool(SOURCE_DB_URL) as source_pool, await asyncpg.create_pool(TARGET_DB_URL) as target_pool, ): builder.provide(SOURCE_DB, source_pool) builder.provide(TARGET_DB, target_pool) yield @coco.fn(memo=True) async def process_record(record: SourceRecord, target_table: postgres.TableTarget[TargetRecord]) -> None: target_table.declare_row(row=TargetRecord( id=record.id, name=record.name.upper(), value=record.value * 2, processed=True, )) @coco.fn async def app_main() -> None: target_table = await postgres.mount_table_target( TARGET_DB, table_name="target_records", table_schema=await postgres.TableSchema.from_class(TargetRecord, primary_key=["id"]), ) source = postgres.PgTableSource( coco.use_context(SOURCE_DB), table_name="source_records", row_type=SourceRecord, ) await coco.mount_each( coco.component_subpath("record"), process_record, source.fetch_rows().items(key=lambda r: r.id), target_table, ) app = coco.App(coco.AppConfig(name="DatabaseTransform"), app_main)

关键点剖析

  • 源表遍历与键提取PgTableSource.fetch_rows()从源表读取记录流,.items(key=lambda r: r.id)将每条记录以r.id为稳定键包装成(key, value)对。mount_each()因此为每条记录建立一个组件,组件子路径为record/{id}——主键即组件身份,记录不变则组件整体跳过(process_record带有memo=True),这是数据库增量 ETL 的核心机制。

  • 双连接池上下文:源库与目标库的连接池都通过@coco.lifespan注册到环境中,async with保证生命周期结束时正确关闭。目标表与源表分属不同库,在同一个 app 内被统一管理。

  • 声明式行目标target_table.declare_row(row=TargetRecord(...))声明"该行应当存在"。引擎对比目标表当前内容,自动生成 insert / update / delete。TableSchema.from_class(TargetRecord, primary_key=["id"])从 dataclass 推断列定义并显式指定主键。

  • memo=True的副作用边界:注意process_record声明了目标行,属于"声明目标状态"而非外部副作用,因此记忆化是安全的——未变化记录不会重复写行。

Pattern 4:LLM 结构化抽取管线

适用场景:调用 LLM 从非结构化文本中抽取结构化数据(主题、实体、摘要等),写入多张目标表。

完整代码示例

import instructor from dataclasses import dataclass from typing import AsyncIterator from pydantic import BaseModel from litellm import acompletion import asyncpg import cocoindex as coco from cocoindex.connectors import postgres DATABASE_URL = "postgres://cocoindex:cocoindex@localhost/cocoindex" PG_DB = coco.ContextKeyasyncpg.Pool _instructor_client = instructor.from_litellm(acompletion, mode=instructor.Mode.JSON) class ExtractedTopic(BaseModel): name: str description: str class ExtractionResult(BaseModel): title: str topics: list[ExtractedTopic] @dataclass class Message: id: int title: str content: str @dataclass class Topic: message_id: int name: str description: str @coco.fn(memo=True) async def extract_and_store( content: str, message_id: int, messages_table: postgres.TableTarget[Message], topics_table: postgres.TableTarget[Topic], ) -> None: result = await _instructor_client.chat.completions.create( model="gpt-4", response_model=ExtractionResult, messages=[{"role": "user", "content": f"Extract topics:\n\n{content}"}], ) messages_table.declare_row(row=Message(id=message_id, title=result.title, content=content)) for topic in result.topics: topics_table.declare_row(row=Topic( message_id=message_id, name=topic.name, description=topic.description, )) @coco.fn async def app_main(input_texts: list[str]) -> None: messages_table = await postgres.mount_table_target( PG_DB, table_name="messages", table_schema=await postgres.TableSchema.from_class(Message, primary_key=["id"]), ) topics_table = await postgres.mount_table_target( PG_DB, table_name="topics", table_schema=await postgres.TableSchema.from_class(Topic, primary_key=["message_id", "name"]), ) for idx, text in enumerate(input_texts): await coco.mount( coco.component_subpath("text", idx), extract_and_store, text, idx, messages_table, topics_table, ) app = coco.App(coco.AppConfig(name="LLMExtraction"), app_main, input_texts=["text1...", "text2..."])

关键点剖析

  • memo=True避免对未变化输入重复调用 LLM:LLM 调用昂贵且有延迟,记忆化在这里价值最高。只要(content, message_id)调用指纹不变、抽取函数代码未变,就完全跳过 LLM 调用,直接复用上次抽取结果。

  • 单个组件声明多张目标表extract_and_store同时向messages_tabletopics_table声明行。CocoIndex 的目标状态机制允许多目标声明,引擎统一 reconcile。这也展示了 1:N 的父子关系如何落库:Topic("message_id", "name")为复合主键。

  • Pydantic 模型约束 LLM 输出:通过instructor.from_litellm(acompletion, mode=instructor.Mode.JSON)构造客户端,response_model=ExtractionResult强制 LLM 返回符合BaseModel结构的 JSON,再逐字段声明为目标行。

  • 显式组件子路径:循环中逐条mount(),并显式传入coco.component_subpath("text", idx)。这与反模式部分强调的"稳定路径"一致——这里idx是输入文本的稳定序号(enumerate(input_texts)),且输入列表本身固定,因此是安全的。

仓库中的 hn_trending_topics 示例 即 LLM 抽取模式的完整落地参考。

Pattern 5:Kafka 流式管线

适用场景:消费 Kafka 消息,实时写入数据库/向量库。

完整代码示例

import json from collections.abc import AsyncIterator from dataclasses import dataclass from confluent_kafka import Message from confluent_kafka.aio import AIOConsumer import cocoindex as coco from cocoindex.connectors import kafka, lancedb LANCE_DB = coco.ContextKeylancedb.LanceAsyncConnection @dataclass class Product: sku: str name: str category: str price: float @coco.lifespan async def coco_lifespan(builder: coco.EnvironmentBuilder) -> AsyncIterator[None]: conn = await lancedb.connect_async("./lancedb_data") builder.provide(LANCE_DB, conn) yield @coco.fn async def process_message(msg: Message, table: lancedb.TableTarget[Product]) -> None: value = msg.value() if value is None: return row = json.loads(value.decode() if isinstance(value, bytes) else value) table.declare_row(row=Product(**{**row, "price": float(row["price"])})) @coco.fn async def app_main() -> None: products_table = await lancedb.mount_table_target( LANCE_DB, table_name="products", table_schema=await lancedb.TableSchema.from_class(Product, primary_key=["sku"]), ) consumer = AIOConsumer({ "bootstrap.servers": "localhost:9092", "group.id": "my-group", "enable.auto.commit": "false", "auto.offset.reset": "earliest", }) items = kafka.topic_as_map(consumer, ["products-topic"]) await coco.mount_each(process_message, items, products_table) app = coco.App(coco.AppConfig(name="KafkaToLanceDB"), app_main)

关键点剖析

  • kafka.topic_as_map()返回LiveMapFeed:在 kafka 源实现 中,topic_as_map(consumer, topics)返回_TopicMapFeedLiveMapFeed[bytes | str, Message]),每条消息以message key为键、完整confluent_kafka.Message为值。要求 consumer未订阅topic_as_map内部处理订阅与分区重平衡回调),且应关闭自动提交(enable.auto.commit: "false"),由 CocoIndex 自动管理 offset。value 为None的 Kafka tombstone 消息会被自动识别为删除事件,也可通过is_deletion谓词自定义删除判定。

  • mount_each()自动识别 live 模式:当itemsLiveMapFeed时,api.py 的mount_each()不会逐项静态挂载,而是内部创建一个_MountEachLiveComponent来持续监听消息流——管线以 live 模式持续运行,新消息到达即被处理,无需手动轮询。process_message中的if value is None: return则是对 tombstone 的额外防御。

  • 运行方式:Kafka 源没有"初始快照",属于纯流式(LiveMapFeed),因此必须以 live 模式运行,例如cocoindex update main.py -L(或app.update_blocking(live=True))。与之对照,walk_dir(live=True)返回的LiveMapView会先扫描当前状态再监听增量,两种模式(catch-up / live)都可用。

仓库中的 csv_to_kafka 示例 展示了数据"写入"Kafka 的方向,rust/examples下还有对应的 csv_to_kafka 与 kafka_consume 的 Rust 版本。

Pattern 6:共享资源上下文管理

适用场景:在多个组件之间共享昂贵的资源(模型、数据库连接、配置)。

完整代码示例

import cocoindex as coco from cocoindex.ops.sentence_transformers import SentenceTransformerEmbedder EMBEDDER = coco.ContextKeySentenceTransformerEmbedder CONFIG = coco.ContextKeydict @coco.lifespan async def coco_lifespan(builder: coco.EnvironmentBuilder) -> AsyncIterator[None]: builder.provide(EMBEDDER, SentenceTransformerEmbedder("all-MiniLM-L6-v2")) builder.provide(CONFIG, {"chunk_size": 1000, "overlap": 200}) yield @coco.fn async def process_item(text: str) -> None: embedder = coco.use_context(EMBEDDER) config = coco.use_context(CONFIG) embedding = await embedder.embed(text) ...

关键点剖析

  • detect_change=True—— 值变化时使 memo 失效:在 context_keys.py 中,ContextKey(key, *, detect_change=False),默认不追踪变化。若某个上下文值影响计算输出(如模型版本、配置参数),应显式设置detect_change=True——值变化时,依赖该上下文的 memo 条目会按上下文指纹自动失效重跑,逻辑类似"代码变更导致 memo 失效"。

  • detect_change=False(默认)—— 不参与指纹的资源:数据库连接、日志器等不影响计算语义的资源,不应开启变化检测,否则连接池重建会无谓地使所有 memo 失效。这正是模式文档给出的默认值理由。

  • ContextKey的 key 是稳定身份,避免跨运行改名ContextKey的字符串 key 同时服务于 memo 指纹、目标状态键(如_TableKey(db_key=db.key, ...))与上下文追踪。改名等于更换身份,会导致旧 memo 全部失效、目标状态被当作"新目标"重建。@coco.lifespan将函数注册到默认 CocoIndex 环境,该环境默认被所有 app 共享。

  • 避免"每个组件加载一次模型":这是最常见的性能反模式(见下节)。模型在 lifespan 中只加载一次,组件内通过coco.use_context(EMBEDDER)获取。

常见反模式与规避方法

模式文档归纳了五类高频反模式,结合源码可以更透彻地理解其危害:

1. 遗漏@coco.fn装饰器

# BAD: Missing decorator async def process_file(file, table): table.declare_row(...) # GOOD: @coco.fn async def process_file(file, table): table.declare_row(...)

@coco.fn是组件处理器(ComponentProcessor)与记忆化的入口。未装饰的函数调用时,目标状态声明不会被引擎追踪,逻辑指纹也不会被记录——管线无法感知其存在。

2. 全量重处理(无记忆化)

# BAD: No memoization @coco.fn # Missing memo=True async def process_file(file, table): embedding = await embedder.embed(await file.read_text()) # Expensive! # GOOD: @coco.fn(memo=True) async def process_file(file, table): ...

没有memo=True时,每次运行都会重新执行昂贵的嵌入/LLM 调用。记忆化的触发条件为"输入指纹 + 逻辑指纹均未变",覆盖输入级(同参数跳过)与代码级(同代码跳过)两层。

3. 组件路径不稳定

# BAD: Using object references or indices for file in files: await coco.mount(coco.component_subpath(file), ...) # Object ref for idx, item in enumerate(items): await coco.mount(coco.component_subpath(idx), ...) # Index changes # GOOD: Use stable identifiers await coco.mount_each(process_file, files.items(), table) # Keys from items()

组件路径是 memo 存储与目标状态管理的身份键。对象引用作为键在跨进程/跨运行序列化时不稳定;数组下标在插入/删除元素后会整体漂移,导致"同名组件内容错位",进而引发错误的增量行为甚至脏数据。正确做法是使用稳定标识:文件的相对路径、记录的主键、Kafka 消息 key,或mount_each()自动使用items()的键。

4. 每个组件重复加载资源

# BAD: Loading model in every component @coco.fn async def process(text): model = SentenceTransformer("model") # Loaded repeatedly! # GOOD: Load once in lifespan, use via context embedder = coco.use_context(EMBEDDER)

模型加载是重型操作(下载权重、初始化 CUDA 上下文),每组件加载一次会拖垮吞吐。应在@coco.lifespan中加载一次并builder.provide(...),组件内通过use_context()获取。

5. 把目标状态与副作用混在一起

# BAD: Side effects not detected @coco.fn async def process(data): requests.post("https://api.example.com", json=data) # Not detected! # GOOD: Only declare target states table.declare_row(row=result)

CocoIndex 的增量正确性建立在"纯函数 + 声明目标状态"之上:引擎通过 reconcile 决定是否重跑,但外部副作用(HTTP 调用、发邮件、写日志文件)是不可观测的。若组件被跳过,副作用不会执行;若组件重跑,副作用可能重复执行——两者都无法保证。因此副作用必须放在目标状态机制之外显式管理,处理函数只负责声明"应该存在什么"。

将模式组合为完整应用

运行管线

上述所有模式的 app 都通过同一套 CLI 驱动(详见 SKILL.md 与 cli.py):

cocoindex update main.py # 运行 app(catch-up 模式) cocoindex update main.py:my_app # 运行指定 app cocoindex update main.py -L # live 模式持续运行 cocoindex update main.py --full-reprocess # 强制全量重处理 cocoindex drop main.py [-f] # 清空并重置全部状态 cocoindex ls [main.py] # 列出 app cocoindex show main.py [--tree] # 查看组件路径树
  • catch-up 模式(默认):每次update扫描全部源、仅处理变化部分、同步目标状态后返回。每次仍需扫描源以发现变化。
  • live 模式:catch-up 结束后持续运行,组件从源持续流式接收变化(文件监听、Kafka 消费),低延迟应用增量。需要两点配合:app 启用 live 模式 + 使用支持 live 的源(LiveMapViewLiveMapFeed)。

模式选型速查

需求选用模式核心 API
文件格式转换 / 清洗Pattern 1walk_dir()+declare_file()
文档分块 + 向量化入库Pattern 2RecursiveSplitter+SentenceTransformerEmbedder+mount_table_target()
库到库 ETLPattern 3PgTableSource.fetch_rows().items(key=...)
LLM 结构化抽取Pattern 4instructor + Pydantic + 多表declare_row()
实时消息处理Pattern 5kafka.topic_as_map()+ live 模式
共享模型 / 连接 / 配置Pattern 6ContextKey+@coco.lifespan+use_context()

版本注意

本文所有代码均为 CocoIndex v1(>=1.0.0)API。v0 时代的@cocoindex.flow_defDataScopecocoindex.sources.*cocoindex.targets.*等符号已在 v1 移除,如遇第三方资料中的旧写法,请一律以本文与 API 参考 为准。

延伸阅读

  • 连接器参考:PostgreSQL、SQLite、LanceDB、Qdrant、SurrealDB、Doris、Kafka 等连接器的完整能力矩阵
  • API 参考:@coco.fnmount*ContextKey等核心 API 速查
  • 项目搭建指南:cocoindex init与依赖配置
  • CocoIndex Skill 总览:核心概念、CLI 命令与 v0/v1 差异对照
  • 可运行示例:files_transform、text_embedding、csv_to_kafka、hn_trending_topics 等(每个示例目录含独立 README)

【免费下载链接】cocoindexIncremental engine for long horizon agents 🌟 Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/15 21:22:06

驾照考试系统源码解析:Java与PHP双后端协作实战

简介:这是一套基于Java、JavaScript、CSS、HTML、PHP等多种语言开发的驾照考试系统完整源码,面向需要学习全栈Web开发或直接部署驾考平台的开发者、学生及项目实践者,也可作为二次开发的基础。资源共300个文件,其中含78个Java源文…

作者头像 李华
网站建设 2026/9/15 21:20:43

在C盘做网站可以吗?老站长揭秘完整流程与风险

在C盘做网站可以吗?老站长揭秘完整流程与风险 很多新手刚接触建站,第一反应就是把项目文件丢进 C 盘。别慌,我见过太多人因为这一招,导致网站上线后频繁报错,备案审核还卡壳。备案流程确实让人一头雾水,但搞懂 完整流程 背后的技术逻辑,你会发现这并非不可逾越的鸿沟。…

作者头像 李华
网站建设 2026/9/15 21:17:50

MATLAB声发射数据分析:滑动窗口计算b值、熵值、CV值等特征

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华