用 CocoIndex 把商品目录变成推荐图谱:LLM 分类提取 + Neo4j 增量构建实战
【免费下载链接】cocoindexIncremental engine for long horizon agents 🌟 Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex
本文以 examples/product_recommendation 示例为核心,讲解如何用 CocoIndex 把一批商品 JSON 自动改造成一张可查询的 Neo4j 推荐图谱:LLM 负责给每个商品打上"它是什么"和"买它还需要什么"两类分类标签,共享的分类节点与两类关系边共同构成"买了 A 的人还需要 B"的推荐引擎,而这一切运行在普通的异步 Python 中。读完本文,你将掌握 CocoIndex 两阶段流水线的写法、@coco.fn(memo=True)的增量重算机制、Neo4j 图目标(节点表 / 关系边)的挂载方式,以及用一条 Cypher 查询输出跨售推荐的方法。
核心思路:推荐信息就藏在商品描述里
一份商品列表里往往隐含着天然的推荐关系——一支钢笔和墨芯、笔记本搭在一起,一台显示器和支架、HDMI 线搭在一起。但这些知识锁在商品描述的自然语言里,人工整理不现实,静态规则又难以覆盖千变万化的商品。示例的思路是:让 LLM 从商品详情文本中提取两类分类标签,再把标签之间的共享关系建成一张图,推荐结论直接从图里查出来,不需要单独训练或部署任何推荐模型。
整个变换在原生 Python 里声明,形如target_state = transformation(source_state);增量处理、变更跟踪、图目标管理等重活由底层的 Rust 引擎承担,因此编辑一个商品只会重新提取一个商品,而不是整份目录。
示例的入口是 examples/product_recommendation/main.py,全文只有一个文件、约 280 行,是理解 CocoIndex 声明式流水线的极佳范本。
图模型:两类节点、两类关系边
推荐图谱由两种节点和两种关系组成:
| 元素 | 类型 | 说明 |
|---|---|---|
Product节点 | 节点 | 每个商品一条(标题、价格),主键为文件名去掉.json |
Taxonomy节点 | 节点 | 每个去重后的分类标签一条(如gel pen、notebook、ink refill),按值作主键,跨商品共享 |
PRODUCT_TAXONOMY边 | 关系 | Product → Taxonomy:这个商品是什么 |
PRODUCT_COMPLEMENTARY_TAXONOMY边 | 关系 | Product → Taxonomy:买这个商品的人可能还需要什么 |
推荐逻辑因此变得非常简洁:某个商品的 complementary 分类,命中另一个商品的 is-a 分类,两者就应当被一起推荐。
在 main.py 中,这两类节点和两类边分别用 dataclass 与neo4j.TableTarget/neo4j.RelationTarget声明:
@dataclass class Product: id: str # 主键——文件名去掉 .json title: str price: float @dataclass class Taxonomy: value: str # 主键——分类标签本身PRODUCT_TAXONOMY与PRODUCT_COMPLEMENTARY_TAXONOMY两条边不带额外负载,连接器直接以(from_id, to_id)推导每条边的主键,因此同一对(产品, 标签)只会出现一条边。
两阶段流水线:为什么必须分两步
由于 Taxonomy 标签在所有商品间共享,示例把流水线拆成两个阶段(自上而下阅读 main.py 即可看到全貌):
- 阶段一(逐商品):每个商品声明自己的
Product节点,调用 LLM 提取分类,并把标签携带给下一阶段; - 阶段二(一次遍历):由唯一一个 graph pass 统一声明去重后的
Taxonomy节点,以及两类关系边。
如果让每个商品各自声明 Taxonomy 节点,同一个gel pen就会被重复建成多个节点;共享节点必须由单一阶段统一拥有,这正是"shared nodes, done right"的关键。
阶段一:LLM 提取 + 节点声明 + 标签携带
@coco.fn(memo=True) # 按内容缓存每次提取——只重跑发生变化的商品 async def extract_taxonomy(detail: str) -> ProductTaxonomyInfo: client = instructor.from_litellm(litellm.acompletion, mode=instructor.Mode.JSON) result = await client.chat.completions.create( model=coco.use_context(LLM_MODEL), response_model=ProductTaxonomyInfo, messages=[{"role": "system", "content": TAXONOMY_PROMPT}, {"role": "user", "content": detail}], ) return ProductTaxonomyInfo.model_validate(result.model_dump()) @coco.fn(memo=True) # 阶段一——逐商品:声明节点、提取、携带标签 async def process_file( file: FileLike, product_table: neo4j.TableTarget[Product], ) -> ProductTaxonomies: raw = json.loads(await file.read_text()) product_id = file.file_path.path.name.removesuffix(".json") price = float(str(raw["price"]).lstrip("$").replace(",", "")) product_table.declare_record(row=Product(id=product_id, title=raw["title"], price=price)) info = await extract_taxonomy(PRODUCT_TEMPLATE.render(**raw)) return ProductTaxonomies( product_id=product_id, taxonomies=[t.name for t in info.taxonomies], complementary=[t.name for t in info.complementary_taxonomies], )值得注意的几个实现细节:
- 价格清洗:
float(str(raw["price"]).lstrip("$").replace(",", ""))能同时处理"$4.99"和"$1,349.00"这类带货币符号与千分位的原始字符串; - 提取入参:商品 JSON 不是直接丢给 LLM,而是先经
PRODUCT_TEMPLATE(Jinja2 模板)渲染成结构化的 Markdown 文本(标题 + Highlights + Description),再作为detail传入; - 结构化解耦:LLM 提取逻辑(
extract_taxonomy)与文件读取逻辑(process_file)是独立的@coco.fn,前者按detail内容 memo 化,后者按文件 memo 化,互不影响缓存粒度。
LLM 提取的响应结构由 Pydantic 模型约束(配合 instructor 的 JSON 模式):
class ProductTaxonomy(pydantic.BaseModel): name: str = pydantic.Field( description=( "A concise noun (or short noun phrase) for the product's core " "functionality, without branding or style. ..." ) ) class ProductTaxonomyInfo(pydantic.BaseModel): taxonomies: list[ProductTaxonomy] = pydantic.Field( description="Taxonomies describing what this product is." ) complementary_taxonomies: list[ProductTaxonomy] = pydantic.Field( description="Taxonomies for complementary products a buyer of this product might also need." )TAXONOMY_PROMPT的系统提示词要求模型"只返回文本中确有支撑的内容",name字段的约束(具体名词、小写、不带品牌与风格修饰、避免office supplies这类过宽分类而倾向pen、printer等具体分类)直接决定了共享标签的质量——标签越规范,跨商品共享命中率越高。
阶段二:一次遍历统一建共享节点与边
@coco.fn # 阶段二——一次遍历拥有共享 Taxonomy 节点 + 两类边 async def build_graph( products: list[ProductTaxonomies], taxonomy_table: neo4j.TableTarget[Taxonomy], product_taxonomy_rel: neo4j.RelationTarget[Any], complementary_rel: neo4j.RelationTarget[Any], ) -> None: labels: set[str] = set() for p in products: labels.update(p.taxonomies) labels.update(p.complementary) for value in labels: taxonomy_table.declare_record(row=Taxonomy(value=value)) for p in products: for t in set(p.taxonomies): product_taxonomy_rel.declare_relation(from_id=p.product_id, to_id=t) for t in set(p.complementary): complementary_rel.declare_relation(from_id=p.product_id, to_id=t)这里的去重逻辑一目了然:先把所有商品的 is-a 与 complementary 标签并成一个set,再逐一声明Taxonomy节点——保证gel pen在全图中只有一个节点,所有商品都指向它。关系边用set()包裹同一商品的标签列表,避免重复边。
环境准备与依赖
示例的依赖定义在 examples/product_recommendation/pyproject.toml:
[project] name = "product-recommendation" version = "0.1.0" description = "CocoIndex example: LLM-extract product taxonomies into a Neo4j recommendation graph." requires-python = ">=3.11" dependencies = [ "cocoindex[neo4j]>=1.0.7", "instructor>=1.0.0", "litellm>=1.0.0", "pydantic>=2.0.0", "jinja2>=3.1.0", ]环境变量模板见 examples/product_recommendation/.env.example:
| 变量 | 默认值 | 说明 |
|---|---|---|
COCOINDEX_DB | ./cocoindex.db | CocoIndex 本地状态库路径(记录增量状态与 memo 缓存) |
OPENAI_API_KEY | (空) | OpenAI 密钥;使用本地模型时可不填 |
LLM_MODEL | openai/gpt-4.1 | LiteLLM 模型标识,可换ollama/llama3.2等任意 provider |
NEO4J_URI | bolt://localhost:7687 | Neo4j Bolt 地址 |
NEO4J_USER | neo4j | Neo4j 用户名 |
NEO4J_PASSWORD | cocoindex | Neo4j 密码 |
按照 README 的四步即可跑通:
1. 启动 Neo4j:
docker run -d -p 7474:7474 -p 7687:7687 -e NEO4J_AUTH=neo4j/cocoindex --name cocoindex-neo4j neo4j:5.26-community2. 配置并安装:
cp .env.example .env # 填入 OPENAI_API_KEY(或改用 LLM_MODEL=ollama/llama3.2) pip install -e .3. 构建图谱—— 示例自带products/目录,包含 9 份示例商品清单(笔、笔记本、显示器配件等):
cocoindex update main在 9 个示例商品上,这条命令会产出9 个Product节点、约 40 个Taxonomy节点,以及两类关系边。
4. 在 Neo4j Browser 中探索推荐—— 打开 http://localhost:7474(neo4j/cocoindex),对图提问。
用一条 Cypher 查询"买了 A 还需要 B"
README 给出了核心的推荐查询:找到所有 is-a 分类命中某支"凝胶笔"互补分类的商品:
// 推荐与任何"gel pen"搭配的商品: // 找到 is-a 分类命中一支笔的 complementary 分类的商品 MATCH (:Taxonomy {value: "gel pen"})<-[:PRODUCT_TAXONOMY]-(:Product) -[:PRODUCT_COMPLEMENTARY_TAXONOMY]->(need:Taxonomy) MATCH (rec:Product)-[:PRODUCT_TAXONOMY]->(need) RETURN DISTINCT rec.title查询分两步走:先沿PRODUCT_COMPLEMENTARY_TAXONOMY从一支笔走到它"可能还需要"的 Taxonomy 集合,再沿PRODUCT_TAXONOMY反向找出所有 is-a 命中这些标签的商品,最后DISTINCT去重。在示例数据上,为一支笔做推荐会得到笔记本(notepad)和多用途纸(multipurpose paper)——正是期望中的跨售组合。图谱本身就是推荐器:没有额外的模型,没有单独的训练流程。
增量更新:编辑一个商品会发生什么
示例最核心的价值在于增量。README 明确了两层机制:
@coco.fn(memo=True)按内容缓存每次 LLM 提取:编辑一个商品的文件内容后,只有该商品的提取会被重跑,其余商品直接命中缓存,不会重新调用 LLM;- 图谱整体做差异(diff)更新:阶段二以"声明期望状态"的方式工作,CocoIndex 会自动新增不再存在的节点/边、删除已没有任何商品支持的节点/边。
两层合起来的效果是:edit one product re-extracts one product, not the catalog——改一个商品只重算一个商品,而不是整份目录。对依赖 LLM 调用(有成本、有延迟)的流水线来说,这意味着后续迭代只付出 O(1) 的增量成本。
源码级原理:ContextKey、挂载与 memo
环境注入:ContextKey与coco.lifespan
连接配置和模型选择都通过上下文键注入,而不是全局变量:
KG_DB = coco.ContextKeyneo4j.ConnectionFactory LLM_MODEL = coco.ContextKeystr @coco.lifespan async def coco_lifespan(builder: coco.EnvironmentBuilder) -> AsyncIterator[None]: builder.provide( KG_DB, neo4j.ConnectionFactory( uri=os.environ.get("NEO4J_URI", "bolt://localhost:7687"), auth=( os.environ.get("NEO4J_USER", "neo4j"), os.environ.get("NEO4J_PASSWORD", "cocoindex"), ), database=os.environ.get("NEO4J_DATABASE", "neo4j"), ), ) builder.provide(LLM_MODEL, os.environ.get("LLM_MODEL", "openai/gpt-4.1")) yielddetect_change=True是"诚实的缓存失效"开关:LLM_MODEL被声明为变更可检测,一旦你在.env里把模型从openai/gpt-4.1换成ollama/llama3.2,流水线会对所有商品针对新模型重新提取,无需手工清缓存。coco.use_context(LLM_MODEL)在extract_taxonomy中读取该上下文。
从源码看,python/cocoindex/connectors/neo4j/_target.py 中的ConnectionFactory持有uri、auth、database三个连接参数,按需创建带认证的异步连接池 driver,并通过query(cypher, params)直接执行单条 Cypher 语句。工厂把 database 名放进连接层(而非表 key),因此同一个 Neo4j 集群上的不同库通过不同的ConnectionFactory/ContextKey对来寻址。
挂载目标:mount_table_target与mount_relation_target
示例在app_main中把四个目标一次性挂载到图数据库:
@coco.fn async def app_main(sourcedir: pathlib.Path) -> None: product_table = await neo4j.mount_table_target( KG_DB, "Product", await neo4j.TableSchema.from_class(Product, primary_key="id"), primary_key="id", ) taxonomy_table = await neo4j.mount_table_target( KG_DB, "Taxonomy", await neo4j.TableSchema.from_class(Taxonomy, primary_key="value"), primary_key="value", ) product_taxonomy_rel = await neo4j.mount_relation_target( KG_DB, "PRODUCT_TAXONOMY", product_table, taxonomy_table ) complementary_rel = await neo4j.mount_relation_target( KG_DB, "PRODUCT_COMPLEMENTARY_TAXONOMY", product_table, taxonomy_table ) ...对照 python/cocoindex/connectors/neo4j/_target.py 的实现:TableTarget.declare_record把行转成字典并校验主键存在,再通过coco.declare_target_state声明该节点的期望状态;RelationTarget.declare_relation则从两端TableTarget取 label 与主键字段,用结构化参数(而非字符串拼接)绑定from_id/to_id,天然避免 Cypher 注入。declare_row = declare_record提供了别名,两者等价。
数据源:localfs.walk_dir
文件来源声明在app_main中:
files = localfs.walk_dir( sourcedir, recursive=True, path_matcher=PatternFilePathMatcher(included_patterns=["**/*.json"]), ) file_coros = [] async for path_key, file in files.items(): file_coros.append( coco.use_mount( coco.component_subpath("file", path_key), process_file, file, product_table, ) ) products: list[ProductTaxonomies] = list(await asyncio.gather(*file_coros))localfs.walk_dir递归遍历products/目录,PatternFilePathMatcher只保留**/*.json;每个文件挂载一个独立的process_file组件(组件子路径按path_key区分),最后asyncio.gather并行收集全部提取结果,作为阶段二的输入。注意这里process_file返回的ProductTaxonomies只是阶段间传递的内部类型,并不直接落到图上。
示例数据一览
examples/product_recommendation/products 下是 9 份结构统一的商品 JSON(每份含title、price、highlights与description),例如 p1.json:
{ "title": "Pilot G2 Premium Gel Roller Pens, Fine Point, 0.7mm, Black Ink, 2/Pack", "price": "$4.99", "highlights": [ "Smooth-writing gel ink for effortless note-taking.", "Comfortable rubber grip for long study sessions.", "Refillable design reduces waste and saves money." ], "description": { "header": "...", "paragraph": "...", "bullets": ["..."] } }9 个商品覆盖了笔(p1)、多用途纸(p2)、橡皮(p3)、螺旋笔记本(p4)、白板(p5)、双肩包(p6)、保温水瓶(p7)、太阳能笔记本(p8)与激光打印机(p9)等品类,正好能演示跨品类互补推荐(文具 + 纸张、设备 + 配件)。换成你自己的商品目录时,只要保持同样的 JSON 结构即可,无需改动任何流水线代码。
小结与扩展方向
回顾这个示例的四个设计要点:
- 共享节点由单一 pass 拥有:
Taxonomy按值去重,gel pen是全图唯一的节点,所有商品共享指向它; - 增量是默认行为:
@coco.fn(memo=True)按内容缓存 LLM 提取,改一个商品只重算一个商品,随后整个图谱做差异更新; - 图即推荐器:无需独立模型,一条 Cypher 沿
complementary → is-a路径即可输出跨售候选; - 纯 Python、你自己的技术栈:提取层是 instructor 叠加 LiteLLM,
LLM_MODEL可切换任意 provider(OpenAI、Ollama 等),没有 DSL。
如果想继续深入,可以在 python/cocoindex/connectors/neo4j/_target.py 中查看节点表、关系边与向量索引的完整目标实现;其他图数据库(FalkorDB、SurrealDB)的连接器也提供了结构一致的TableTarget/RelationTarget/mount_relation_targetAPI(分别见 python/cocoindex/connectors/falkordb/_target.py 与 python/cocoindex/connectors/surrealdb/_target.py),同一套声明式写法可以直接迁移。把商品目录变成推荐图谱,本质上就是把"藏在散文里的知识"结构化——增量引擎保证这个过程可以随目录一起长期演进。
【免费下载链接】cocoindexIncremental engine for long horizon agents 🌟 Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考