使用Snowpipe时如何实现数据upsert操作?
背景
我非常喜欢使用Snowpipe,但在使用过程中无法套用现有的upsert逻辑。
我目前的upsert逻辑如下:
create temp table temp_table (like target); copy into temp_table from @snowflake_stage; begin transaction; delete from target using temp_table where target.pk = temp_table.pk; insert into target select * from temp_table; end transaction; drop table temp_table;
但Snowpipe仅允许定义单条copy命令,无法执行多命令序列。
已尝试方案
我曾考虑过使用Tasks & Streams,但Tasks似乎不支持事务(单任务内无法执行多条查询)。我也尝试过使用MERGE语法,但该语法要求显式指定要INSERT的列。
例如我无法执行如下无需指定插入列的操作:
merge into src using temp_table on src.pk = temp_table.pk when not matched then insert;
请问还有其他可以在使用Snowpipe的同时完成数据upsert的方案吗?
解决方案
中转表+流+存储过程+任务组合方案
首先纠正一个认知误区:单条Task本身确实不支持直接写多条SQL语句,但Task可以调用存储过程,存储过程内完全支持多语句事务,完全兼容你原有的upsert逻辑。
具体实现步骤:- 新建一张永久中转表,结构和目标表完全一致,用于承接Snowpipe同步的数据
- Snowpipe的COPY目标直接设置为这张中转表,不需要修改Snowpipe的基础配置逻辑
- 为中转表创建变更流(Stream),用于捕获中转表内的新增数据
- 编写存储过程,在存储过程内实现你原有的事务性upsert逻辑,处理完成后清空中转表数据即可
- 创建调度任务,触发条件设置为流中有新增数据时自动执行,任务逻辑直接调用上述存储过程即可。
动态生成MERGE语句方案
如果你想简化链路,不想维护额外的流和任务,可以通过动态SQL解决MERGE需要手动指定列的问题。
你可以在存储过程中先从INFORMATION_SCHEMA.COLUMNS查询出目标表的全量列名,自动拼接成MERGE语句中需要的INSERT列列表和对应VALUE值,不需要手动维护列名,表结构变更后也能自动适配。动态表(Dynamic Table)方案
如果你使用的Snowflake版本支持动态表特性,可以采用更轻量化的方案:将Snowpipe的落地表作为源表,定义动态表的刷新逻辑为你需要的upsert逻辑,Snowflake会自动按你设置的刷新频率做增量同步,不需要手动管理流、任务和事务。
内容的提问来源于stack exchange,提问作者Reynold Rogers

