CocoIndex 增量写入 Snowflake表目标Table Target实战与源码解析【免费下载链接】cocoindexIncremental engine for long horizon agents Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex导读本文以仓库中的 Snowflake 目标示例 为主线完整讲解如何用 CocoIndex 把数据流中声明的行rows增量地写入 SnowflakeCocoIndex 负责目标状态管理——按需创建数据库、Schema 与表用 SnowflakeMERGE实现按主键的 upsert并在数据流不再声明某行时自动DELETE删除。读完本文你将掌握cocoindex update的完整运行流程、.env各配置项的含义、TableTarget的声明方式以及连接器底层的两级目标状态表级 DDL 行级 DML实现原理。Snowflake 作为 CocoIndex 表目标职责划分示例的核心思路是Snowflake 只做表存储CocoIndex 全权负责目标状态。具体来说CocoIndex在需要时创建 database、schema 和 table用 SnowflakeMERGE写入行删除数据流中不再声明的行。Snowflake只负责保存由 CocoIndex 管理的表数据。源数据留在 Python 里示例把三笔订单记录放在 Python 内存中这样无需搭建外部数据源就能直观观察 Snowflake 目标的真实行为——声明三笔订单、计算order_total、写入 Snowflake最后用 Snowflake Python connector 把表读回来验证。默认创建的对象默认情况下该数据流会写入以下对象对象默认值创建方WarehouseCOCOINDEX_DEMO_WH你运行示例之前手动创建DatabaseCOCOINDEX_DEMO_DBCocoIndexSchemaPUBLICCocoIndex如需要TableCOCOINDEX_ORDERSCocoIndex关键前提Warehouse 必须在数据流运行前已存在。Snowflake 连接器只会创建 database、schema 和 table不会创建 warehouse——这是 Snowflake 的约束DDL/DML 必须挂在一个已存在的 warehouse 上执行。数据流工作方式原文档将整个流程概括为四步对应 main.py 中的代码coco_lifespan读取SNOWFLAKE_*环境变量向 CocoIndex 提供 Snowflake 连接配置mount_table_target声明 Snowflake 表目标及其主键process_order把每条源订单转换成写入 Snowflake 的行形状计算order_totalcocoindex update main将声明的行与表内容对账reconcile。前置条件与安装以下命令都应在示例目录内执行cd examples/snowflake_target1. 安装依赖从源码 checkout 安装示例的依赖声明见 pyproject.toml其中cocoindex[snowflake]1.0.7会一并安装snowflake-connector-pythonpip install -e ../..[snowflake] -e .2. 复制环境模板并填写.envcp .env.example .envmain.py通过load_dotenv()加载.env。各环境变量说明如下取自原文档变量必填如何选择COCOINDEX_DB是本地 CocoIndex 状态数据库。示例用默认的./cocoindex.db即可。SNOWFLAKE_ACCOUNT是Snowflake 账户标识例如ORGNAME-ACCOUNTNAME。不要带.snowflakecomputing.com后缀。SNOWFLAKE_USER是Snowflake 登录名。SNOWFLAKE_PASSWORD是SNOWFLAKE_USER的密码。SNOWFLAKE_ROLE否会话使用的角色。留空则使用用户的默认角色。SNOWFLAKE_WAREHOUSE是用于执行 DDL 和 DML 的已存在 warehouse。README 使用COCOINDEX_DEMO_WH。SNOWFLAKE_DATABASE是目标数据库名。角色有权限时 CocoIndex 会创建它。SNOWFLAKE_SCHEMA是目标 Schema 名。默认是PUBLIC。SNOWFLAKE_TABLE是目标表名。默认是COCOINDEX_ORDERS。对应到 main.pyDATABASE、SCHEMA、TABLE_NAME三个常量直接读取这些环境变量未设置时回退到默认值。你也可以在 Snowflake worksheet 中查询当前会话值用于核对.envSELECT CURRENT_ORGANIZATION_NAME() AS organization_name, CURRENT_ACCOUNT_NAME() AS account_name, CURRENT_USER() AS user_name, CURRENT_ROLE() AS role_name, CURRENT_WAREHOUSE() AS warehouse_name;对于SNOWFLAKE_ACCOUNT使用 Snowflake Account Details 中显示的账户标识较新的账户通常形如ORGANIZATION_NAME-ACCOUNT_NAME。3. 在 Snowflake 中创建演示 warehouseCREATE WAREHOUSE IF NOT EXISTS COCOINDEX_DEMO_WH WAREHOUSE_SIZE XSMALL AUTO_SUSPEND 60 AUTO_RESUME TRUE;如果使用其他 warehouse 名称记得同步更新.env中的SNOWFLAKE_WAREHOUSE。4. 角色权限要求试用账户用ACCOUNTADMIN即可。若使用更窄的角色至少需要warehouse 的USAGE权限创建或使用目标 database 的权限创建或使用目标 schema 的权限创建目标表以及执行MERGE和DELETE的权限。认证方式示例使用用户名 密码认证对应 coco_lifespan 中构造的snowflake.ConnectionConfigsnowflake.ConnectionConfig( accountos.environ[SNOWFLAKE_ACCOUNT], useros.environ[SNOWFLAKE_USER], passwordos.environ[SNOWFLAKE_PASSWORD], warehouseos.environ.get(SNOWFLAKE_WAREHOUSE), roleos.environ.get(SNOWFLAKE_ROLE) or None, )在源码 python/cocoindex/connectors/snowflake/_target.py 中ConnectionConfig是一个 frozen dataclass字段即account、user、password以及可选的warehouse、role连接时通过snowflake.connector.connect(**kwargs)建立见 源码 _connect。注意密钥对key-pair认证不在本示例覆盖范围内当前连接配置接受的是account、user、password、warehouse、role五个字段实际验证走的是密码路径。核心代码拆解声明表目标与行示例的数据模型在 main.py源侧SourceOrderorder_id、customer、product、quantity、unit_price、status、attributes字典目标侧SnowflakeOrder在源模型基础上多出order_total字段。三笔示例订单SAMPLE_ORDERS保存在内存中见 main.py。声明表目标app_main中通过mount_table_target声明目标main.pytable await snowflake.mount_table_target( SNOWFLAKE, table_nameTABLE_NAME, table_schemaawait snowflake.TableSchema.from_class( SnowflakeOrder, primary_key[order_id], ), databaseDATABASE, schemaSCHEMA, )这里TableSchema.from_class从SnowflakeOrder这个 dataclass 推导出全部列定义并指定order_id为主键。在源码 python/cocoindex/connectors/snowflake/_target.py 中mount_table_target先通过table_target()构造TargetState校验表名/库名/schema 名/列名均为合法标识符再await coco.mount_target(...)挂载返回可直接使用的TableTarget。声明行process_order是一个带memoTrue的函数通过table.declare_row声明目标行main.pycoco.fn(memoTrue) async def process_order(order, table): table.declare_row( rowSnowflakeOrder( order_idorder.order_id, ... order_totalround(order.quantity * order.unit_price, 2), ... ) )mount_each将三笔订单逐一喂给process_ordermain.py。在底层 TableTarget.declare_row 中行会被转换为dict取出主键值作为目标状态的 key然后调用coco.declare_target_state声明该行应存在。源码级原理两级目标状态与 SQL 生成Snowflake 连接器在 python/cocoindex/connectors/snowflake/_target.py 中实现了一套两级目标状态系统表级Table level创建/删除 Snowflake 表DDL行级Row level在表内 upsert/删除行DML。表级 handler建库、建 schema、建表、增量改列_TableHandler源码负责 DDL。当表需要创建时它会依次执行可在测试 python/tests/connectors/test_snowflake_target.py 中看到精确的 SQL 断言CREATE DATABASE IF NOT EXISTS COCOINDEX_DEMO_DB CREATE SCHEMA IF NOT EXISTS COCOINDEX_DEMO_DB.PUBLIC CREATE TABLE COCOINDEX_DEMO_DB.PUBLIC.COCOINDEX_ORDERS (order_id VARCHAR NOT NULL, ..., PRIMARY KEY (order_id))值得注意的实现细节主键列自动NOT NULL_column_sql中主键列即使声明为可空也会被强制NOT NULL见 源码 _column_sql。增量 Schema 演进表级跟踪记录由CompositeTrackingRecord组成主键签名 每列的类型/可空性子记录见 源码 _table_composite_tracking_record_from_spec。当声明的行结构发生变化时_apply_column_actions会执行ALTER TABLE ... ADD COLUMN / DROP COLUMN IF EXISTS等增量 DDL源码。managed_by语义target.ManagedBy源码 python/cocoindex/connectorkits/target.py区分SYSTEMCocoIndex 管理生命周期与USER用户管理。表级对账通过statediff.resolve_system_transition源码 python/cocoindex/connectorkits/statediff.py只对 system 管理的部分做差异计算避免误删用户自己建的表。行级 handlerMERGE 与 DELETE_RowHandler源码 _apply_actions把 CocoIndex 对账产出的行级动作翻译成 SQLupsert 走MERGE按主键匹配匹配则更新非主键列不匹配则插入。SQL 由_merge_sql生成源码测试 python/tests/connectors/test_snowflake_target.py 验证了其形态MERGE INTO DB.SCHEMA.TABLE AS target USING (SELECT %s AS col1, PARSE_JSON(%s) AS col2, ...) AS source ON target.pk source.pk WHEN MATCHED THEN UPDATE SET col2 source.col2 WHEN NOT MATCHED THEN INSERT (col1, col2, ...) VALUES (source.col1, source.col2, ...)删除走DELETE单列主键用WHERE pk IN (...)复合主键用(pk1 %s AND pk2 %s) OR ...展开见 源码 _delete_sql 与对应测试 test_delete_sql_supports_single_and_composite_primary_keys。VARIANT 列用PARSE_JSONPython 中的dict/list会被编码为 JSON 字符串并在 SQL 中以PARSE_JSON(%s)还原为 SnowflakeVARIANT使attributes:channel::string这类 JSON 路径查询可以直接使用源码 _source_select_sql。事务语义同一批动作在一个连接内执行成功后commit异常时rollback并重抛源码。Python 类型到 Snowflake 类型的默认映射TableSchema.from_class依据字段的 Python 类型推断 Snowflake 列类型源码 _LEAF_TYPE_MAPPINGSPython 类型Snowflake 类型备注boolBOOLEANintNUMBERfloatFLOATdecimal.DecimalNUMBERstrVARCHARbytesBINARYuuid.UUIDVARCHAR编码为字符串datetime.dateDATEdatetime.timeTIMEdatetime.datetimeTIMESTAMP_TZdatetime.timedeltaNUMBER编码为总秒数dict/list/ 嵌套记录 / 联合 /AnyVARIANT使用PARSE_JSON类型映射的正确性由测试 test_table_schema_maps_python_types_to_snowflake_types 直接断言。若需覆盖默认映射可用typing.Annotatedsnowflake.SnowflakeType指定列类型和自定义编码器例如Annotated[int, snowflake.SnowflakeType(NUMBER(38, 0))]见 源码 SnowflakeType 与测试 test_snowflake_type_override_is_used。运行与预期输出构建/更新目标表cocoindex update main首次运行会创建 database、schema、表并用MERGE写入三行。之后每次运行只做增量对账。读取 Snowflake 中的行python main.pymain.py底部的print_rows()main.py直接用snowflake.connector连接并把表读回打印其中attributes:channel::string演示了从 VARIANT 中提取字段(ORD-1001, Summit Labs, mechanical keyboard, 2, 259.0, paid, web) (ORD-1002, Beacon Retail, standing desk, 1, 399.0, paid, sales) (ORD-1003, Ridgeview Health, noise cancelling headphones, 3, 599.97, pending, partner)注意order_total是quantity * unit_price四舍五入到两位的结果2 × 129.50 259.0、1 × 399.00 399.0、3 × 199.99 599.97。体验增量更新编辑 main.py 中的SAMPLE_ORDERS然后再次运行cocoindex update main python main.py行为如下修改某行只有发生变化的订单行被 upsertMERGE的WHEN MATCHED THEN UPDATE分支删除某行如果从SAMPLE_ORDERS中移除某笔订单下一次运行时 CocoIndex 会把对应行从 Snowflake 中DELETE掉。这正是 CocoIndex 增量引擎的核心价值CocoIndex 的本地状态库默认./cocoindex.db记录每一行的指纹行级 handler 用fingerprint_object计算跟踪记录见 源码 reconcile只有指纹变化或缺失的行才会触发真实写入避免了全表重建。验证与测试连接器行为在仓库测试 python/tests/connectors/test_snowflake_target.py 中有完整覆盖单元级非法标识符拒绝、类型映射、SnowflakeType覆盖、MERGE/DELETESQL 生成、行编码的确定性VARIANT 序列化按键排序、表级 DDL 执行序列实时级test_live_snowflake_upsert_and_delete测试源码在配置了SNOWFLAKE_*环境变量时才会运行它会创建一张随机表名COCOINDEX_TEST_UUID的真实表写入、更新、删除并断言最终只剩一行且值为更新后的内容最后清理表——验证了 upsert/delete 在真实 Snowflake 上的语义。如果你在本地配置了 Snowflake 凭据可以仿照该测试的思路用临时表名 自动清理的方式安全地验证本文描述的全部行为。【免费下载链接】cocoindexIncremental engine for long horizon agents Star if you like it!项目地址: https://gitcode.com/GitHub_Trending/co/cocoindex创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考