You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在ADF中实现Snowflake迁移项目的批量文件记录数校验

解决方案:ADF中实现数据文件与控制文件的记录数校验

一、整体流程设计

核心思路为分层遍历+双元数据获取+动态校验,具体流程如下:

  1. 遍历XYZ容器下的Batch1至Batch10所有文件夹
  2. 在每个Batch文件夹内,筛选并遍历所有.ctl控制文件
  3. 对每个控制文件,匹配对应的.csv.gz数据文件
  4. 提取控制文件记录数,统计数据文件实际有效记录数(减去表头1行)
  5. 校验两者是否一致,一致则触发Snowflake数据加载流程

二、具体步骤实现

1. 获取Batch文件夹列表

  • 新增Get Metadata活动,数据源指向ADLS的XYZ容器,Field list选择Child items,获取所有Batch文件夹的列表。
  • 将该活动结果传入第一个For Each活动,遍历每个Batch文件夹。

2. 筛选并遍历Batch内的控制文件

在第一个For Each活动内部:

  • 新增第二个Get Metadata活动,数据源指向当前遍历的Batch文件夹,Field list选择Child items。
  • 添加Filter活动,对返回的文件列表进行过滤,只保留.ctl后缀的文件,过滤条件为:
    @endsWith(item().name, '.ctl')
    
  • 将过滤后的控制文件列表传入第二个For Each活动,遍历每个控制文件。

3. 解析控制文件的记录数与对应数据文件名

在第二个For Each活动内部:

  • 新增Lookup活动,读取当前控制文件内容(ctl为单行文本,需将First row only设为true)。
  • 用ADF表达式解析ctl的单行数据,提取关键信息:
    # 获取对应的数据文件名(假设ctl格式为"文件名,记录数,日期")
    @split(activity('Lookup_CTL').output.firstRow.Column1, ',')[0]
    
    # 获取控制文件记录数并转为整数
    @int(split(activity('Lookup_CTL').output.firstRow.Column1, ',')[1])
    

4. 统计数据文件的实际有效记录数

ADF原生活动无法直接读取压缩文件行数,提供两种可行方案:

方案A:Azure Function计算行数(推荐大文件场景)

  • 创建Azure Function,配置ADLS存储连接权限,编写代码读取gzip压缩文件,跳过表头后统计行数,核心Python示例:
    import gzip
    import os
    from azure.storage.blob import BlobServiceClient
    
    def main(file_path: str) -> int:
        blob_service_client = BlobServiceClient.from_connection_string(os.environ['AZURE_STORAGE_CONNECTION_STRING'])
        blob_client = blob_service_client.get_blob_client(container='XYZ', blob=file_path)
        
        with gzip.open(blob_client.download_blob(), 'rt') as f:
            lines = list(f)
            return len(lines) - 1  # 减去表头行
    
  • 在ADF中新增Azure Function活动,传入当前数据文件的路径,获取返回的实际记录数。

方案B:临时导入统计行数(轻量小文件场景)

  • 创建Snowflake临时表或Azure SQL临时表,用Copy Activity将当前csv.gz文件导入,设置Skip header rows为1。
  • 新增Lookup活动,执行查询获取实际行数:
    SELECT COUNT(*) AS actual_count FROM temp_table
    

5. 校验与加载触发

  • 新增If Condition活动,判断控制文件记录数与实际记录数是否一致:
    @equals(variables('ctl_record_count'), variables('actual_record_count'))
    
  • 条件为真时,执行Copy Activity将数据文件加载到Snowflake;条件为假时,可记录错误日志(如写入ADLS错误文件夹)。

三、关键注意事项

  • 文件匹配逻辑:确保ctl与csv.gz文件名严格对应,可通过表达式替换后缀实现:@replace(item().name, '.ctl', '.csv.gz')
  • 性能优化:若文件数量大,可将For Each活动的Sequential设为false开启并行处理,注意控制ADF并发上限。

内容的提问来源于stack exchange,提问作者Deepika

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 09:01:40