如何用Azure Data Factory Dataflow从分离数组读取行列生成CSV
问题描述
API返回的响应结构中,列定义和行数据分别存储在独立数组中,示例响应如下:
[ { "FrameType": "DataTable", "TableId": 1, "TableKind": "PrimaryResult", "TableName": "PrimaryResult", "Columns": [ { "ColumnName": "RecordID", "ColumnType": "string" }, { "ColumnName": "RecordNumber", "ColumnType": "integer" }, { "ColumnName": "IsValid", "ColumnType": "boolean" }, { "ColumnName": "Remarks", "ColumnType": "string" } ], "Rows": [ [ "record1", 123, true, "test1" ], [ "record2", 456, false, "test2" ] ] } ]
需要通过Azure Data Factory将该响应转换为CSV文件存储到Blob中,之前尝试使用派生列和Flatten转换未成功,原Data Flow JSON如下:
{ "name": "df_ADX_REST_TagsInfo", "properties": { "type": "MappingDataFlow", "typeProperties": { "sources": [ { "linkedService": { "referenceName": "ls_ADX_REST_API", "type": "LinkedServiceReference" }, "name": "ADXRestApi" } ], "sinks": [ { "linkedService": { "referenceName": "linkservice_cloudshellstorageforbh", "type": "LinkedServiceReference" }, "name": "WriteColumnstoCSV" } ], "transformations": [ { "name": "FilterAPIResponse" }, { "name": "flatten1" }, { "name": "derivedColumn1" }, { "name": "derivedColumn2" }, { "name": "select1" } ], "scriptLines": [ "parameters{", " access_token as string", "}", "source(output(", " body as (Cancelled as boolean, Columns as (ColumnName as string, ColumnType as string)[], FrameType as string, HasErrors as boolean, IsProgressive as boolean, Rows as string[][], TableId as short, TableKind as string, TableName as string, Version as string),", " headers as [string,string]", " ),", " allowSchemaDrift: true,", " validateSchema: false,", " format: 'rest',", " timeout: 30,", " requestInterval: 0,", " entity: '/v2/rest/query',", " headers: ['Content-Type' -> 'application/json; charset=utf-8', 'Authorization' -> ($access_token), 'Accept' -> 'application/json', 'Host' -> 'help.kusto.windows.net'],", " httpMethod: 'POST',", " body: ('{\"db\":\"bhTestDB\", \"csl\":\"TagsInfo | summarize arg_max(ingestion_time(),*) by TimeseriesId\"}'),", " paginationRules: ['supportRFC5988' -> 'true'],", " responseFormat: ['type' -> 'json', 'documentForm' -> 'documentPerLine']) ~> ADXRestApi", "ADXRestApi filter(body.TableKind=='PrimaryResult') ~> FilterAPIResponse", "FilterAPIResponse foldDown(unroll(body.Columns),", " mapColumn(", " Columns = body.Columns.ColumnName", " ),", " skipDuplicateMapInputs: false,", " skipDuplicateMapOutputs: false) ~> flatten1", "FilterAPIResponse derive(Tags = toString(unfold(split(toString(body.Rows), '],[')))) ~> derivedColumn1", "derivedColumn1 derive(Tags = replace(replace(Tags,'[[',''),']]','')) ~> derivedColumn2", "derivedColumn2 select(mapColumn(", " Tags", " ),", " skipDuplicateMapInputs: true,", " skipDuplicateMapOutputs: true) ~> select1", "flatten1 sink(allowSchemaDrift: true,", " validateSchema: false,", " input(", " {{} as string", " ),", " format: 'delimited',", " container: 'data',", " columnDelimiter: ',',", " escapeChar: '\\\\',", " quoteChar: '\\\"',", " columnNamesAsHeader: true,", " skipDuplicateMapInputs: true,", " skipDuplicateMapOutputs: true,", " saveOrder: 1) ~> WriteColumnstoCSV" ] } } }
可行实现方案
核心思路
先展开行数据数组,再将列名与行数据逐一映射,最终导出为结构化CSV。
修正后的Data Flow脚本
{ "name": "df_ADX_REST_TagsInfo", "properties": { "type": "MappingDataFlow", "typeProperties": { "sources": [ { "linkedService": { "referenceName": "ls_ADX_REST_API", "type": "LinkedServiceReference" }, "name": "ADXRestApi" } ], "sinks": [ { "linkedService": { "referenceName": "linkservice_cloudshellstorageforbh", "type": "LinkedServiceReference" }, "name": "WriteColumnstoCSV" } ], "transformations": [ { "name": "FilterAPIResponse" }, { "name": "FlattenRows" }, { "name": "MapColumns" } ], "scriptLines": [ "parameters{", " access_token as string", "}", "source(output(", " body as (Cancelled as boolean, Columns as (ColumnName as string, ColumnType as string)[], FrameType as string, HasErrors as boolean, IsProgressive as boolean, Rows as array[], TableId as short, TableKind as string, TableName as string, Version as string),", " headers as [string,string]", " ),", " allowSchemaDrift: true,", " validateSchema: false,", " format: 'rest',", " timeout: 30,", " requestInterval: 0,", " entity: '/v2/rest/query',", " headers: ['Content-Type' -> 'application/json; charset=utf-8', 'Authorization' -> ($access_token), 'Accept' -> 'application/json', 'Host' -> 'help.kusto.windows.net'],", " httpMethod: 'POST',", " body: ('{\"db\":\"bhTestDB\", \"csl\":\"TagsInfo | summarize arg_max(ingestion_time(),*) by TimeseriesId\"}'),", " paginationRules: ['supportRFC5988' -> 'true'],", " responseFormat: ['type' -> 'json', 'documentForm' -> 'documentPerLine']) ~> ADXRestApi", "ADXRestApi filter(body.TableKind=='PrimaryResult') ~> FilterAPIResponse", // 展开Rows数组,将每行转为单独记录 "FilterAPIResponse flatten(unroll(body.Rows),", " mapColumn(", " RowData = body.Rows", " ),", " skipDuplicateMapInputs: false,", " skipDuplicateMapOutputs: false) ~> FlattenRows", // 映射列名与行数据 "FlattenRows derive(", " RecordID = RowData[1],", " RecordNumber = RowData[2],", " IsValid = RowData[3],", " Remarks = RowData[4]", ") ~> MapColumns", // 导出到Blob CSV "MapColumns sink(allowSchemaDrift: true,", " validateSchema: false,", " format: 'delimited',", " container: 'data',", " columnDelimiter: ',',", " escapeChar: '\\\\',", " quoteChar: '\\\"',", " columnNamesAsHeader: true,", " skipDuplicateMapInputs: true,", " skipDuplicateMapOutputs: true,", " saveOrder: 1) ~> WriteColumnstoCSV" ] } } }
关键步骤说明
- 过滤目标表:保留
TableKind为PrimaryResult的数据,确保只处理需要的表结构。 - 展开行数据:通过Flatten转换的
unroll(body.Rows)将嵌套的行数组展开,每个子数组成为一条独立记录,存储到RowData字段。 - 映射列与值:用派生列转换,根据
Columns中的列名,从RowData数组按索引提取对应值(注意Data Flow数组索引从1开始),生成结构化列。 - 导出CSV:将映射后的结构化数据写入Blob存储,启用
columnNamesAsHeader确保CSV包含表头。
动态列名适配(可选)
如果列名可能动态变化,可使用动态列映射:
- 在FlattenRows后添加派生列,使用表达式
map(body.Columns, (col, index) => col.ColumnName -> RowData[index+1])自动关联列名和行数据。 - 需保持Data Flow的
allowSchemaDrift: true和validateSchema: false配置,支持动态列结构。
内容的提问来源于stack exchange,提问作者Pathrudu
相关产品推荐
相关产品推荐

