如何通过Data Factory将Azure Blob存储数据增量加载到Azure SQL数据库
可行实现方案
方案1:Azure Data Factory 原生水印增量加载(适配已有DF全量加载流程,改动最小)
- 前置要求:你的JSON文件记录需要有唯一标识列(比如业务ID、流水号)和时间戳列(比如记录创建时间、更新时间),没有的话可直接用JSON文件的Blob上传时间作为辅助判断依据
- 操作步骤:
- 先在Azure SQL中创建一个水印表,存储每次加载完成的最大时间戳/最大ID,初始值设置为上次全量加载时对应的最大时间/最大ID
- 在现有DF管道中新增查找活动,先读取水印表中的当前最大水印值
- 调整复制活动的源端配置:连接Blob存储的JSON文件,新增筛选条件,仅读取大于查找活动返回的最大水印值的记录,源端筛选可直接用JSONPath表达式,示例写法:
$.record_create_time > @activity('查水印').output.firstRow.max_watermark - 复制活动的目标端保持原有Azure SQL表配置,写入模式选择追加
- 管道末尾新增更新活动,把本次加载拿到的最大水印值回写到水印表,供下次调度使用
- 把管道触发频率设置为每周,和JSON文件更新时间对齐即可
- 兜底校验配置:如果担心出现重复数据,可在复制活动前新增存储过程活动,执行
MERGE语句,匹配到唯一ID的旧记录不插入,仅插入未匹配到的新记录,示例SQL片段:MERGE INTO 目标业务表 t USING 增量临时表 s ON t.唯一标识ID = s.唯一标识ID WHEN NOT MATCHED THEN INSERT (列1,列2,列3) VALUES (s.列1,s.列2,s.列3);
方案2:Blob文件属性判断增量(适合JSON无明确时间戳列的场景)
- 操作步骤:
- 直接用Blob自带的
LastModified属性,或者提前开启Blob存储的版本控制能力 - DF管道启动后先获取上次加载时间,再列出目标Blob路径下所有
LastModified/版本时间大于上次加载时间的JSON文件 - 把筛选出的新增JSON文件全量读入SQL临时表,再用
MERGE语句和正式业务表比对,插入未匹配的新记录 - 最后把本次加载时间更新到水印表即可
- 直接用Blob自带的
方案3:轻量无管道方案(适合数据量不大的场景)
- 操作步骤:
- 用Azure Functions配置Blob触发规则,每周新JSON上传后自动触发函数执行
- 函数逻辑内读取JSON内容,同时连接Azure SQL拉取现有记录的唯一ID列表
- 过滤出JSON中不在现有ID列表的新记录,批量插入到SQL表即可
内容的提问来源于stack exchange,提问作者Olgaraa
相关产品推荐
相关产品推荐

