如何通过Azure Data Factory自动将新增日期文件夹文件加载至现有表
创建Azure Pipeline自动加载分层存储的新增文件到现有表
针对你的年/月/日层级存储结构(如2023/01/01),每日新增日期文件夹并自动加载文件到现有表的需求,以下提供两种主流实现方案:
方案一:基于Azure DevOps Pipeline(通用CI/CD自动化)
前置准备
- 确认存储账户(如ADLS Gen2)的层级结构合规,每个日期文件夹下的4个文件字段与目标表匹配
- 拥有Azure DevOps项目权限,以及存储账户、目标数据库(如Azure SQL DB/Synapse)的读写权限
- 准备数据加载脚本(PowerShell/Python均可)
步骤1:配置服务连接
在Azure DevOps项目的项目设置→服务连接中,创建两类连接:
- Azure资源管理器连接:用于统一访问存储账户和目标数据库
- 或分别创建存储账户服务连接和Azure SQL数据库服务连接,确保权限覆盖读取存储文件、写入目标表
步骤2:编写数据加载脚本
以PowerShell为例,核心逻辑是对比已加载日志找出新增文件夹,批量加载文件并记录日志:
# 配置参数(建议用Pipeline变量替换硬编码值) $storageAccountName = "$(StorageAccountName)" $containerName = "$(ContainerName)" $targetServer = "$(TargetSqlServer)" $targetDB = "$(TargetDatabase)" $targetTable = "$(TargetTableName)" # 从日志表获取最后加载日期 $lastLoadedDate = (Invoke-SqlCmd -ServerInstance $targetServer -Database $targetDB -Query "SELECT ISNULL(MAX(LoadedDate), '2000-01-01') FROM LoadLog").Column1 # 获取指定年月下的所有日期文件夹(可动态替换年月参数) $rootDir = "$(Year)/$(Month)" $folders = Get-AzDataLakeGen2ChildItem -AccountName $storageAccountName -FileSystem $containerName -Directory $rootDir | Where-Object { $_.Name -match "\d{4}/\d{2}/\d{2}" } # 筛选未加载的日期文件夹 $newFolders = $folders | Where-Object { [DateTime]$_.Name.Split('/')[-1] -gt [DateTime]$lastLoadedDate } foreach ($folder in $newFolders) { $folderPath = $folder.Name # 遍历文件夹下的4个文件 $files = Get-AzDataLakeGen2ChildItem -AccountName $storageAccountName -FileSystem $containerName -Directory $folderPath foreach ($file in $files) { $fileUri = "https://$storageAccountName.dfs.core.windows.net/$containerName/$folderPath/$($file.Name)" # 使用bcp命令加载数据(根据文件格式调整参数) bcp "$targetDB.dbo.$targetTable" in "$fileUri" -S "$targetServer" -U "$(SqlUsername)" -P "$(SqlPassword)" -c -t "," } # 更新加载日志 $loadDate = $folder.Name.Split('/')[-1] Invoke-SqlCmd -ServerInstance $targetServer -Database $targetDB -Query "INSERT INTO LoadLog (LoadedDate, LoadTime) VALUES ('$loadDate', GETDATE())" }
步骤3:配置Pipeline任务
- 在Azure DevOps中新建Pipeline,选择代码源(如Azure Repos Git)并导入脚本
- 添加Azure PowerShell任务:
- 关联已创建的Azure服务连接
- 设置脚本路径为仓库中的脚本文件
- 在变量选项卡配置所有参数(如
StorageAccountName、SqlUsername等),敏感值设为机密变量
- 可选添加前置检查任务:验证存储账户和数据库连接是否正常
步骤4:设置自动触发
- 日程触发:在Pipeline触发器中选择每日执行(如凌晨2点),确保前一日文件已全部上传
- 事件触发:通过Azure Event Grid监听存储账户的文件夹创建事件,自动触发Pipeline:
- 在Azure门户为存储账户创建Event Grid订阅,事件类型选择“Blob创建”
- 终点选择Azure DevOps Pipeline的触发端点
- 在Pipeline中启用事件触发器并关联该订阅
步骤5:错误处理与通知
- 添加发送电子邮件任务,设置条件为“任务失败时执行”,通知相关人员
- 启用Pipeline的失败重试机制,设置重试次数(如2次)
方案二:基于Azure Data Factory(ADF)Pipeline(数据集成专用)
步骤1:创建ADF资源
- 在Azure门户新建Azure Data Factory实例,配置存储账户(ADLS Gen2)和目标数据库的链接服务
步骤2:配置数据集
- 源数据集:选择存储账户,设置文件路径为
@concat(pipeline().parameters.Year, '/', pipeline().parameters.Month, '/', pipeline().parameters.Day),启用“递归”以读取文件夹下所有文件 - 目标数据集:选择目标数据库,指定现有表名
步骤3:构建Pipeline
- 添加Get Metadata活动:获取指定年月下的所有日期文件夹路径
- 添加Filter活动:对比ADF的运行日志或自定义日志表,筛选出未加载的文件夹
- 添加For Each活动:遍历筛选后的文件夹,执行以下子活动:
- Copy Data活动:将源数据集的文件数据复制到目标表
- Stored Procedure活动:调用数据库存储过程,记录本次加载的日期和状态
步骤4:设置触发机制
- 日程触发器:设置每日触发,动态传递当日的年、月、日参数
- 事件触发器:通过Event Grid监听存储账户的文件夹创建事件,触发Pipeline并传递文件夹路径参数
步骤5:监控与告警
- 在ADF的监控面板查看Pipeline运行状态
- 配置Azure Monitor告警,当Pipeline失败时发送通知
内容的提问来源于stack exchange,提问作者user21510438
相关产品推荐
相关产品推荐

