如何将Azure Synapse每日更新的Delta表UPSERT至Azure Data Explorer?
Azure Synapse Delta表到ADX的每日UPSERT同步方案
核心思路
Delta表通过事务日志管理增量文件,直接全量替换会浪费资源,最优方案是先获取当日新增/更新的Delta文件,再在ADX中执行UPSERT操作,既保证数据一致性又提升同步效率。
具体实现步骤
1. 获取Delta表的增量文件路径
在Synapse中通过Spark作业或Delta Lake API,读取Delta表的事务日志,筛选出上次同步后新增的parquet文件:
from delta.tables import DeltaTable # 加载目标Delta表 delta_table = DeltaTable.forPath(spark, "abfss://gold@dev.dfs.core.windows.net/db/goldtestdb/") # 筛选上次同步时间之后的事务记录,提取新增文件路径 history_df = delta_table.history().filter("timestamp > '2024-01-01 00:00:00'") # 替换为实际上次同步时间 file_list = [] for row in history_df.collect(): if row.operation == "WRITE": # 从operationParameters中提取写入的文件路径前缀,拼接完整路径 file_prefix = row.operationParameters.get("target", "") file_list.extend([f"{file_prefix}/{f}" for f in row.operationMetrics.get("addedFiles", "").split(",")]) # 整理为ADX ingest支持的格式(单引号包裹,多文件用逗号分隔) adx_file_paths = ",".join([f"'abfss://gold@dev.dfs.core.windows.net/db/goldtestdb/{f}'" for f in file_list])
也可以用Synapse管道的Lookup活动调用上述逻辑,将文件路径作为参数传递给后续的ADX Command活动。
2. ADX侧执行UPSERT操作
ADX的.merge命令是实现UPSERT的核心,推荐先将增量数据导入临时表,再与目标表合并(性能更优):
// 步骤1:将增量文件导入临时表 .set-or-replace #DeltaTempTable <| ingest into table #DeltaTempTable ({adx_file_paths};impersonate) with (format='parquet') // 步骤2:执行UPSERT到目标表 .merge into YourTargetTable as target using #DeltaTempTable as source on target.YourPrimaryKey == source.YourPrimaryKey // 替换为实际主键列 when matched then update set * when not matched then insert *
如果不需要临时表,也可以直接通过external_data()读取ADLS文件完成合并(适合文件量少的场景):
.merge into YourTargetTable as target using external_data( CityId:string, CityName:string, // 按实际表结构定义列名和类型 h@'abfss://gold@dev.dfs.core.windows.net/db/goldtestdb/*/*.parquet;impersonate' ) with (format='parquet') as source on target.CityId == source.CityId when matched then update set * when not matched then insert *
3. 自动化每日同步
将上述步骤整合到Synapse管道:
- 配置Spark作业活动执行增量文件路径获取逻辑
- 将文件路径作为参数传入Azure Data Explorer Command活动,执行Kusto命令
- 设置管道为每日定时触发,确保ADX数据每日自动同步更新
注意事项
- 权限:确保ADX服务主体拥有ADLS存储容器的读取权限,
impersonate参数需符合当前身份验证配置 - 主键:必须明确目标表的主键列,否则
.merge无法正确匹配更新/插入数据 - 增量筛选:务必准确设置上次同步时间,避免重复同步或遗漏数据
内容的提问来源于stack exchange,提问作者LJRB
相关产品推荐
相关产品推荐

