如何使用KQL查询实现ADX到Azure存储的增量加载
仅加载ADX最新数据的增量同步方案
要实现每周从ADX增量加载最新数据到Azure存储,核心是通过时间标识字段过滤出上次同步后新增/更新的数据,以下是具体实现步骤:
关键前提
确保你的ADX表具备以下任一字段:
IngestionTime:ADX自动生成的字段,记录数据被摄入到集群的时间(默认开启,若未开启需在表创建时指定.ingestion_time())- 业务自定义时间字段:比如
LastUpdatedAt,记录数据最后更新的业务时间
实现步骤
1. 记录上次同步的时间戳
需要维护一个"上次同步完成时间"的标记,可存储在:
- Azure存储的元数据文件(比如JSON格式)
- Azure SQL Database的专门记录表
- Spark作业的参数或配置中(适合固定周期调度,比如每周同步上周数据)
2. 编写增量KQL查询
根据使用的时间字段,编写过滤最新数据的KQL:
示例1:使用ADX自带的IngestionTime
假设上次同步时间是2024-05-01T00:00:00Z,查询如下:
YourTableName | where IngestionTime > datetime(2024-05-01T00:00:00Z) | project 字段1, 字段2, ... # 按需指定要加载的字段
示例2:使用业务自定义时间字段
如果用LastUpdatedAt作为更新标识:
YourTableName | where LastUpdatedAt > datetime(2024-05-01T00:00:00Z) | project 字段1, 字段2, ...
3. 修改Spark连接器代码
将增量查询整合到Spark代码中,并动态传入上次同步时间:
# 假设从存储读取上次同步时间,实际需替换为读取元数据的逻辑 last_sync_time = "2024-05-01T00:00:00Z" query = f""" YourTableName | where IngestionTime > datetime({last_sync_time}) | project Column1, Column2, Column3 # 替换为你的实际字段 """ dataframe = spark.read. \ format("com.microsoft.kusto.spark.datasource"). \ option("kustoCluster", kustoOptions["kustoCluster"]). \ option("kustoDatabase", kustoOptions["kustoDatabase"]). \ option("kustoQuery", query). \ option("kustoAadAppId", kustoOptions["kustoAadAppId"]). \ option("kustoAadAppSecret", kustoOptions["kustoAadAppSecret"]). \ option("kustoAadAuthorityID", kustoOptions["kustoAadAuthorityID"]). \ load() # 导出数据到Azure存储(示例为ADLS Gen2) dataframe.write. \ format("parquet"). \ mode("append"). # 用追加模式避免覆盖已有数据 option("path", "abfss://container@storageaccount.dfs.core.windows.net/your-path"). \ save() # 更新上次同步时间为当前时间,需写入到之前维护的元数据存储位置 from datetime import datetime new_sync_time = datetime.utcnow().strftime("%Y-%m-%dT%H:%M:%SZ") # 这里添加写入new_sync_time到元数据存储的代码,比如写入Azure Blob
额外注意事项
- 如果ADX表存在更新操作(不只是新增),仅用IngestionTime可能会遗漏已更新的历史数据,此时必须依赖业务自定义的更新时间字段。
- 调度作业时,确保每周同步的时间窗口覆盖上周的所有更新,比如每周一同步上周日之前的所有新数据。
- 首次同步时,可将
last_sync_time设为一个很早的时间(比如1970-01-01T00:00:00Z),加载全量数据后再记录当前时间作为下次同步的基准。
内容的提问来源于stack exchange,提问作者krish
相关产品推荐
相关产品推荐

