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

配置文件驱动的ADF多源数据湖Raw到Staging层动态Ingestion方案咨询

ADF动态加载Raw到Staging层解决方案(配置文件驱动)

一、配置文件设计(核心)

在DataLake中创建结构化配置文件(推荐JSON格式,也可使用CSV),统一管理所有待加载的源路径、文件类型、日期规则等。示例配置如下:

[
  {
    "source_container": "raw-container",
    "source_folder": "sales/region-1",
    "target_staging_table": "staging_db.dbo.sales",
    "file_type": "csv",
    "date_filter_config": {
      "filter_type": "relative", // 可选absolute/relative
      "relative_days": -7, // 加载近7天数据
      "date_column": "transaction_date" // 文件内容中的日期字段(用于内容过滤)
    },
    "file_naming_pattern": "sales_*.csv" // 可选:文件名匹配通配符
  },
  {
    "source_container": "raw-container",
    "source_folder": "customers",
    "target_staging_table": "staging_db.dbo.customers",
    "file_type": "json",
    "date_filter_config": {
      "filter_type": "absolute",
      "start_date": "2024-01-01",
      "end_date": "2024-01-31",
      "date_column": "signup_date"
    },
    "file_naming_pattern": "*.json"
  }
]

二、ADF流水线核心架构

采用配置驱动+循环遍历模式,核心组件流程如下:

  1. Lookup活动:读取上述配置文件,将配置数据以数组形式传入后续流程。

    • 数据源选择存储账户中的配置文件,开启"返回所有行"选项。
  2. For Each活动:遍历Lookup返回的每个配置项,可根据资源情况调整并行处理度。

    • 迭代项设置为@activity('Lookup_Config').output.value。
  3. 分支逻辑(If Condition):在For Each内部,根据@item().file_type判断文件类型,分别处理CSV/JSON:

    • CSV处理分支:
      • 使用Copy Data活动,源数据集参数化:
        • 容器路径绑定@item().source_container
        • 文件夹路径绑定@item().source_folder
        • 文件格式选择CSV,开启"首行作为标题"
        • 内容过滤:若需基于日期字段过滤,在源的"查询"中写入动态SQL(适用于带Schema的CSV):
          SELECT * FROM @{dataset().source_path} WHERE transaction_date BETWEEN '@{item().date_filter_config.start_date}' AND '@{item().date_filter_config.end_date}'
          
        • 目标数据集指向Staging层的SQL表/ADLS文件夹,自动或手动配置字段映射。
    • JSON处理分支:
      • 同样使用Copy Data活动,源数据集选择JSON格式:
        • 配置JSON读取模式(单行/多行数组)
        • 路径绑定逻辑同CSV,内容过滤可通过JSON路径表达式实现,或前置Data Flow活动做过滤转换。

三、日期范围过滤的两种实现方式

  • 文件名级过滤:如果Raw层文件按日期命名(如sales_20240201.csv),直接在Copy Data的源路径中使用通配符,结合配置的日期范围动态生成:
    @concat(item().source_folder, '/sales_', formatDateTime(adddays(utcnow(), item().date_filter_config.relative_days), 'yyyyMMdd'), '*.csv')
    
  • 内容级过滤:如果日期字段在文件内容中,通过Copy Data的查询(CSV)或Data Flow的过滤转换(JSON)实现,确保只加载指定日期范围内的数据。

四、参数化与扩展性设计

  • 流水线参数:添加config_file_path参数,支持不同环境(dev/test/prod)使用不同配置文件。
  • 数据集参数:将容器、文件夹路径、文件类型设为参数,避免重复创建多个数据集。
  • 扩展支持更多文件类型:只需在配置中新增file_type(如parquet),并在For Each中添加对应的分支处理逻辑。

五、初期测试建议

  1. 单独测试单个配置项:创建仅加载某一个CSV文件夹的流水线,验证路径、过滤、目标写入是否正确。
  2. 验证日期过滤逻辑:分别测试绝对日期和相对日期两种场景,检查加载的数据范围是否符合预期。
  3. 测试错误处理:故意配置错误路径/文件类型,验证Catch活动是否能捕获异常并记录日志。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 15:02:35