通过ADF数据流处理Blob转SQLDB时映射错误求助
Azure Data Factory 数据流映射错误修复方案
问题说明
参考相关方案在Azure Data Factory(ADF)中通过数据流实现「获取最后修改的文件并将JSON转CSV」的流程时,出现映射错误导致流程失败。相关数据流代码如下:
source(output( {$schema} as string, type as string, items as (type as string, properties as (columns as (type as string, items as (type as string)[]), rows as (type as string, items as (type as string, items as (type as string)[])[])), required as string[])[] ), allowSchemaDrift: true, validateSchema: false, ignoreNoFilesFound: false, purgeFiles: true, modifiedAfter: (toTimestamp(1679961600000L)), documentForm: 'arrayOfDocuments') ~> source1 source1 select(mapColumn( columns ), skipDuplicateMapInputs: true, skipDuplicateMapOutputs: true) ~> select1 source1 derive(rows = unfold(rows)) ~> derivedColumn1 derivedColumn1 derive(rows = replace(replace(toString(rows),'[',''),']','')) ~> derivedColumn2 select1 derive(columns = replace(replace(toString(columns),'[',''),']','')) ~> derivedColumn3 derivedColumn2 rank(asc(rows, true), caseInsensitive: true, output(id as long), dense: true) ~> rank1 rank1 select(mapColumn( rows, id ), skipDuplicateMapInputs: true, skipDuplicateMapOutputs: true) ~> select2 select2, select3 union(byName: false)~> union1 union1 sort(asc(id, true)) ~> sort1 derivedColumn3 aggregate(groupBy(columns), column = first(columns)) ~> aggregate1 aggregate1 select(mapColumn( column ), skipDuplicateMapInputs: true, skipDuplicateMapOutputs: true) ~> select3 sort1 select(mapColumn( rows ), skipDuplicateMapInputs: true, skipDuplicateMapOutputs: true) ~> select4 select4 sink(allowSchemaDrift: true, validateSchema: false, skipDuplicateMapInputs: true, skipDuplicateMapOutputs: true) ~> DS2csv
错误根源分析
- 节点执行顺序混乱:
union1调用select3时,select3还未被定义(select3在union1之后才生成),ADF数据流按拓扑顺序执行,会直接报找不到对象的错误。 - 列路径映射错误:
columns和rows字段实际嵌套在items.properties下,直接映射columns/rows会找不到对应字段。 - 嵌套数组展开逻辑错误:原代码直接对
rows执行unfold,但rows是二维嵌套数组,未处理内层数组会导致数据格式混乱。 - Union操作逻辑错误:
select2(含rows、id)和select3(含column)结构完全不匹配,按位置合并(byName: false)后无法生成符合CSV要求的结构。
修复步骤及修正后代码
1. 调整节点执行顺序
确保select3的生成节点(derivedColumn3、aggregate1)在union1之前执行,保证union操作时依赖的节点已存在。
2. 修正列映射路径
将columns和rows的路径修正为items.properties.columns.items和items.properties.rows,确保能正确读取嵌套字段。
3. 修正嵌套数组展开逻辑
先展开外层rows数组,再展开内层rows.items数组,确保每一行数据被正确解析。
4. 修正Union与输出逻辑
给表头加标识ID,按ID排序保证表头在前,合并后统一输出为CSV的行内容。
修正后的完整数据流代码
source(output( {$schema} as string, type as string, items as (type as string, properties as (columns as (type as string, items as (type as string)[]), rows as (type as string, items as (type as string, items as (type as string)[])[])), required as string[])[] ), allowSchemaDrift: true, validateSchema: false, ignoreNoFilesFound: false, purgeFiles: true, modifiedAfter: (toTimestamp(1679961600000L)), documentForm: 'arrayOfDocuments') ~> source1 # 处理表头列 source1 select(mapColumn( columns = items.properties.columns.items ), skipDuplicateMapInputs: true, skipDuplicateMapOutputs: true) ~> select1 select1 derive(columns_str = replace(replace(toString(columns),'[',''),']','')) ~> derivedColumn3 derivedColumn3 aggregate(groupBy(columns_str), header = first(columns_str)) ~> aggregate1 aggregate1 derive(id = 0L) ~> select3 # 给表头加标识ID,确保排序后在最前 # 处理数据行 source1 derive(rows = unfold(items.properties.rows)) ~> derivedColumn1 derivedColumn1 derive(rows_items = unfold(rows.items)) ~> derivedColumn1_1 derivedColumn1_1 derive(row_str = replace(replace(toString(rows_items),'[',''),']','')) ~> derivedColumn2 derivedColumn2 rank(asc(row_str, true), caseInsensitive: true, output(id as long), dense: true) ~> rank1 rank1 select(mapColumn( row_str, id ), skipDuplicateMapInputs: true, skipDuplicateMapOutputs: true) ~> select2 # 合并表头和数据行,按ID排序 select3, select2 union(byName: true) ~> union1 union1 sort(asc(id, true)) ~> sort1 sort1 select(mapColumn( content = iif(id == 0, header, row_str) ), skipDuplicateMapInputs: true, skipDuplicateMapOutputs: true) ~> select4 # 输出到CSV select4 sink(allowSchemaDrift: true, validateSchema: false, skipDuplicateMapInputs: true, skipDuplicateMapOutputs: true, fileExtension: '.csv', quoteAllText: true) ~> DS2csv
内容的提问来源于stack exchange,提问作者Jeff
相关产品推荐
相关产品推荐

