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

通过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

错误根源分析

  1. 节点执行顺序混乱:union1调用select3时,select3还未被定义(select3在union1之后才生成),ADF数据流按拓扑顺序执行,会直接报找不到对象的错误。
  2. 列路径映射错误:columns和rows字段实际嵌套在items.properties下,直接映射columns/rows会找不到对应字段。
  3. 嵌套数组展开逻辑错误:原代码直接对rows执行unfold,但rows是二维嵌套数组,未处理内层数组会导致数据格式混乱。
  4. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 06:15:23