使用 CocoIndex Rust SDK 构建 PostgreSQL 数据源增量索引管道从 pgvector 存储到语义搜索【免费下载链接】cocoindexIncremental engine for long horizon agents Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex本指南以 CocoIndex 仓库中的 examples/rust/postgres_source 为例完整讲解如何用 CocoIndex 的 Rust SDK 实现一条Postgres 源表 → 逐行计算派生字段与文本嵌入 → 写入带 pgvector 索引的输出表 → 语义相似度查询的增量索引管道。读完本文你将掌握postgres::read_table读取源表、#[cocoindex::function(memo)]记忆化计算、postgres::mount_table_target声明式目标表与declare_vector_index向量索引声明、以及 pgvector 余弦距离检索的完整 Rust 写法并理解其增量更新的底层机制。示例概览该示例是 Python 版postgres_source示例的 Rust 移植管线职责完全对齐从 Postgres源表source_products读取商品行为每行计算派生字段total_value price × amount并基于拼接后的商品描述生成 384 维文本嵌入将结果写入另一张 Postgres输出表带 pgvector 索引输出由 CocoIndex 的声明式TableTarget托管最后对输出表执行 pgvector 相似度搜索。所有核心代码集中在一个文件 src/main.rs 中依赖声明见 Cargo.toml建表与种子数据见 prepare_source_data.sql。与 Python 示例的逐项对照原 README 用一张对照表清晰标出了 Rust 移植与 Python 原版的对应关系这里逐项展开说明关注点PythonRust本示例源数据行postgres.PgTableSource(...).fetch_rows()postgres::read_table::SourceProduct(db, source_products)逐行计算coco.fn(memoTrue) process_product#[cocoindex::function(memo)] process_product输出存储postgres.TableTarget 向量索引postgres::mount_table_targetdeclare_vector_index嵌入模型sentence-transformers/all-MiniLM-L6-v2fastembed的AllMiniLML6V2同一模型384 维对应地Python 版的完整实现位于 examples/postgres_source/main.py其中使用asyncpg连接池、pgvector.asyncpg注册向量类型并通过coco.fn(memoTrue)装饰逐行计算函数。Rust 版则把同样的结构映射到 CocoIndex Rust SDK 的postgres模块与属性宏体系。嵌入模型细节两端都使用sentence-transformers/all-MiniLM-L6-v2输出维度固定为 384。Rust 侧在 src/main.rs 中通过常量EMBED_MODEL与EMBED_DIM 384声明并通过SentenceTransformerEmbedder::load(EMBED_MODEL)见 src/main.rs加载模型输出表 schema 中嵌入列的类型即vector(384)见 src/main.rs。模型加载依赖通过 Cargo feature 开启见 Cargo.tomlcocoindex { path ../../../rust/sdk/cocoindex, features [postgres, fastembed] }。增量语义memo 跳过、变更重算、删除对账原 README 用三条陈述概括了该管道的增量行为这三条也正是 CocoIndex 增量引擎在源表 派生目标表场景下的核心保证未变化的源行被 memo 跳过重复运行cargo run -- index时内容未变的行不会重新计算嵌入。其机制是#[cocoindex::function(memo)]按行内容对逐行计算进行记忆化——详见 src/main.rs注释明确指出memoized by the row content, so unchanged rows skip the embedding work on re-runs。变化的源行被重算并更新输出行源表中被UPDATE过的行其派生输出行会被 upsert 更新。被删除的源行自动对账移除源表中被DELETE的行其派生输出行会被托管式TableTarget自动删除不会留下孤儿数据。从底层实现看这三条保证由两套机制协同完成源码见 rust/sdk/cocoindex/src/postgres.rsMemo 记忆化process_product是 memo 化的逐行函数源行内容即 memo 键未变行直接命中缓存目标状态对账reconcileTableTarget会把本次声明的行与上一轮跟踪的行做差量对账——变更的行 upsert、未变的行跳过、不再声明的行删除。模块文档对此有明确描述changed rows are upserted, unchanged rows are skipped, and rows no longer declared are deleted见 postgres.rs。此外read_table在实现上会在一次REPEATABLE READ, READ ONLY事务内读取整张表保证全表观察为一致快照见 postgres.rs。因此每次UPDATE/DELETE/INSERT后重跑cargo run -- index即可观察增量重处理效果。环境准备与数据初始化1. 启动带 pgvector 的 Postgres示例要求一个已安装pgvector扩展的 Postgres 实例并将连接串通过POSTGRES_URL环境变量指向它export POSTGRES_URLpostgres://cocoindex:cocoindexlocalhost/cocoindex export SOURCE_DATABASE_URL$POSTGRES_URL # 可选缺省时与 POSTGRES_URL 相同SOURCE_DATABASE_URL用于把源库指向另一个数据库例如源表与输出表分库的场景在 src/main.rs 中两个 URL 的解析逻辑如下POSTGRES_URL缺省时回退到postgres://cocoindex:cocoindexlocalhost/cocoindexSOURCE_DATABASE_URL缺省时直接复用POSTGRES_URL。dotenvy依赖见 Cargo.toml允许把这两个变量写进仓库目录下的.env文件main()中的dotenvy::dotenv().ok()会自动加载见 src/main.rs。2. 建表与种子数据psql $SOURCE_DATABASE_URL -f prepare_source_data.sqlprepare_source_data.sql 完成三件事删除并重建源表source_products其复合主键为(product_category, product_name)见 prepare_source_data.sql并带modified_time时间戳列插入 5 行种子商品数据覆盖 Electronics无线耳机、智能手机、笔记本、Appliances咖啡机、Sports跑鞋三个品类modified_time以NOW() - INTERVAL ...制造时间梯度见 prepare_source_data.sql使用ON CONFLICT (product_category, product_name) DO NOTHING保证重复执行幂等。索引管道实现详解数据模型SourceProduct对应源表一行未声明的多余列会被忽略OutputProduct是写入输出表的行结构比源行多出total_value与embedding两个派生字段见 src/main.rs#[derive(Clone, Serialize, Deserialize)] struct SourceProduct { product_category: String, product_name: String, description: String, price: f64, amount: i64, } #[derive(Clone, Serialize, Deserialize)] struct OutputProduct { product_category: String, product_name: String, description: String, price: f64, amount: i64, total_value: f64, embedding: Vecf32, }逐行计算memo 化process_product是整条管道唯一的计算函数拼接出用于嵌入的完整描述、计算total_value、调用嵌入器得到 384 维向量见 src/main.rs#[cocoindex::function] async fn process_product(ctx: Ctx, product: SourceProduct) - ResultOutputProduct { let full_description format!( Category: {}\nName: {}\n\n{}, product.product_category, product.product_name, product.description ); let total_value product.price * product.amount as f64; let embedding ctx.get_key(EMBEDDER)?.embed(full_description).await?; Ok(OutputProduct { /* ... */ }) }其中EMBEDDER、DB、SOURCE_DB均为ContextKey通过Environment::builder().provide_key(...)注入见 src/main.rs。ContextKey以名字 状态指纹参与环境识别嵌入模型名或数据库连接变化时会触发对应环境的变更检测。输出表 Schema 与向量索引声明output_schema()用postgres::TableSchema::new声明输出表的列与主键见 src/main.rs嵌入列类型直接写vector(384)。主键必须与源表一致地使用(product_category, product_name)——这正是增量对账时用于 upsert/删除的关键。app_main中先挂载目标表并声明 HNSW 向量索引见 src/main.rslet target postgres::mount_table_target(ctx, DB, TABLE, output_schema()?, Some(PG_SCHEMA)).await?; target.declare_vector_index( ctx, embedding, postgres::VectorIndexOptions { method: hnsw, ..Default::default() }, )?;VectorIndexOptions的默认值与可选字段定义在 postgres.rs字段默认值说明nameNone索引名缺省时使用列名metriccosine距离度量支持cosine/l2/ippgvector opclass 映射见 postgres.rsmethodivfflat索引方法本示例显式改为hnswlistsNone仅ivfflat使用的lists参数mNone仅hnsw使用的m参数ef_constructionNone仅hnsw使用的ef_construction参数vector_index_with_clause在生成 DDL 时按索引方法门控参数lists只在ivfflat下输出m/ef_construction只在hnsw下输出避免生成非法 DDL见 postgres.rs。向量索引是目标表的附件attachment会被创建/重建/删除以匹配声明状态见 postgres.rs。读取源表并挂载逐行计算let source_db ctx.get_key(SOURCE_DB)?; let products: VecSourceProduct postgres::read_table(source_db, source_products).await?; let outputs mount_each!( products .into_iter() .map(|p| (format!({}|{}, p.product_category, p.product_name), p)), |product| process_product(ctx, product) ).await?;postgres::read_table将源表整体读为VecTT: Deserialize未知列忽略见 postgres.rsmount_each!为每行生成以category|name为稳定键的任务并并发执行process_product。最后遍历结果逐行声明到目标表见 src/main.rsfor output in outputs { target.declare_row(ctx, output)?; }declare_row内部按主键计算稳定键声明该行的期望状态见 postgres.rs行状态在 reconcile 时以廉价指纹跟踪未变化的行直接跳过不落盘每列数据见 postgres.rs。目标表 DDL 由系统托管define_table会自动CREATE SCHEMA IF NOT EXISTS、启用CREATE EXTENSION IF NOT EXISTS vector、按声明 schema 建表并对照information_schema做列增删对账见 postgres.rs。运行索引与查询按 README 给出的完整流程# 1. 建表 种子数据 psql $SOURCE_DATABASE_URL -f prepare_source_data.sql # 2. 读源表 - 嵌入 - 写输出表重复运行即增量 cargo run -- index # 3. 对输出做语义搜索 cargo run -- query wireless headphonesmain()按第一个参数分派命令见 src/main.rsindex缺省命令连接目标库与源库、加载嵌入器、构建Environment状态库落在CARGO_MANIFEST_DIR/.cocoindex_db、运行app_main并打印统计信息query 文本用sqlx直连数据库执行 pgvector 相似度查询与目标写入分离sqlx依赖见 Cargo.toml。查询实现query()先用同一嵌入器对查询文本编码再执行余弦距离排序的 SQL见 src/main.rslet rows sqlx::query(format!( SELECT product_category, product_name, description, amount, total_value, \ embedding $1::vector AS distance \ FROM \{PG_SCHEMA}\.\{TABLE}\ ORDER BY distance ASC LIMIT $2 )) .bind(vec_lit) .bind(TOP_K) .fetch_all(pool)即 pgvector 的余弦距离算子distance越小越相似输出时用1.0 - distance换算成相似度分数并限制返回前TOP_K 5条见 src/main.rs。同样的查询逻辑在 Python 版中以asyncpg 原生 SQL 呈现见 examples/postgres_source/main.py两端的 SQL 形状完全一致。验证增量按 README 的指引对source_products执行数据变更后重跑索引命令即可观察增量行为-- 例更新一行触发重算 upsert UPDATE source_products SET price 179.99 WHERE product_name Wireless Headphones; -- 例删除一行触发输出行自动对账删除 DELETE FROM source_products WHERE product_name Coffee Maker; -- 例插入一行触发新行首次处理 INSERT INTO source_products (product_category, product_name, description, price, amount) VALUES (Sports, Yoga Mat, Non-slip yoga mat for home workouts, 39.99, 100);随后重跑cargo run -- index未触及的行被 memo 跳过不重新计算嵌入变更/新增行被处理并写入被删除行的输出行由TableTarget自动对账移除。可以用psql对比输出表coco_examples_v1.output的行数与内容来验证三种行为。源码路径速查示例 READMEexamples/rust/postgres_source/README.md示例主程序examples/rust/postgres_source/src/main.rs建表与种子数据examples/rust/postgres_source/prepare_source_data.sql依赖声明examples/rust/postgres_source/Cargo.tomlPython 对应版本examples/postgres_source/main.py 与 examples/postgres_source/README.mdSDK 中 Postgres 源/目标实现rust/sdk/cocoindex/src/postgres.rsread_table/read_table_items、TableSchema/ColumnDef、TableTarget/mount_table_target、VectorIndexOptions、表与行级 reconcile 逻辑均在此文件【免费下载链接】cocoindexIncremental engine for long horizon agents Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考