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

Azure DataBricks输出无法在ADF ForEach循环中遍历的问题

问题:Databricks Notebook输出无法传入ADF ForEach循环

执行ADF的ForEach循环时失败,报错:

模板操作'MainForEach1'执行失败:'foreach'表达式'@activity('test_notebook').output.runOutput'的计算结果类型为'String',结果必须是有效的数组。

我的Databricks代码如下:

final_df = spark.sql("select * from table1")

def call(row):
    row = row.asDict()
    payload = dict()
    payload["id"] = row.get("t1")
    payload["name"] = "test"
    payload["name2"] = row.get("cname")
    payload["add"] = row.get("add1")
    payload["add2"] = row.get("add2")
    payload["c1"] = row.get("c1")
    payload["s1"] = row.get("s1")
    payload["country"] = row.get("Country")
    payload["code"] = row.get("code")
    response = requests.post(api_endpoint,json=payload,headers=headers)
    return response

list_data=list()
# testing purpose
x = final_df.collect()[1:5]
for row in x:
    api_response =call(row)
    temp = {"details":{},"reason":""}
    temp["edetails"]["input"] = row.asDict()
    temp["edetail"]["output"] = api.text
    if api.status_code != 200:
        temp["s"] = "f"
        temp["reason"] = api.text
    else:
        temp["s"] = "s"
    list_data.append(temp)

dbutils.notebook.exit(list_data)

当前Notebook生成的输出如下,但ADF不接受:

[{'EDetails': {'Input': {'test': '3231232', 'cname': 'sdasa', 'add1': None, 'add2': None, 'c1': None, 's1': None, 'c2': None, 'code': None}, 'Output': '{"Message":"An error has occurred."}'}, 'Reason': '{"Message":"An error has occurred."}', 's': 'f'}, {'EDetails': {'Input': {'test': '24431212', 'cname': 'somename', 'add1': 'sfasasd', 'add2': 'dsadasdasd', 'c1': 'swqwewqe', 's1': ' ', 'c1': 'asddadasd', 'code': 'sads2323'}, 'Output': '{"test":"sad12321312312"}'}, 'Reason': '', 's': 's'}, {'EDetails': {'Input': {'test': '23232323', 'cname': "ffffff", 'add1': '43ddsdsd', 'add2': 'NA', 'c1': 'sdfdsfsdf', 's': ' ', 'c1': 'Us', 'code': '323dsds'}, 'Output': '{"test":"0232323"}'}, 'Failure': '', 's1': 's'}, {'EDetails': {'Input': {'test': '33sdsd', 'cname': "sdasdasdas", 'add1': 'sdadsad', 'add2': 'NA', 'c1': 'sadsd', 's1': ' ', 'c1': 'km', 'code': '34eerer'}, 'Output': '{"test":"sadsadsads"}'}, 'Failure': '', 's': 's'}]


解决方案

问题本质

直接用dbutils.notebook.exit(list_data)输出Python列表时,Databricks会把列表转成Python原生字符串格式(比如用单引号、保留None值),这不是标准JSON数组,ADF会将其识别为字符串类型,不符合ForEach要求的数组输入格式。

修复步骤

  1. 引入json模块,把Python列表序列化为标准JSON字符串,确保格式符合ADF解析要求
  2. 处理Python的None值,自动转换为JSON的null
  3. 修正代码里的拼写错误(比如temp["edetails"]和temp["edetail"]键名不一致、api_response与api变量混用的问题)

修改后的代码

import json
import requests

final_df = spark.sql("select * from table1")

def call(row):
    row = row.asDict()
    payload = dict()
    payload["id"] = row.get("t1")
    payload["name"] = "test"
    payload["name2"] = row.get("cname")
    payload["add"] = row.get("add1")
    payload["add2"] = row.get("add2")
    payload["c1"] = row.get("c1")
    payload["s1"] = row.get("s1")
    payload["country"] = row.get("Country")
    payload["code"] = row.get("code")
    response = requests.post(api_endpoint, json=payload, headers=headers)
    return response

list_data = list()
# testing purpose
x = final_df.collect()[1:5]
for row in x:
    api_response = call(row)
    temp = {"EDetails": {}, "Reason": "", "s": ""}
    # 统一键名避免拼写错误
    temp["EDetails"]["Input"] = row.asDict()
    temp["EDetails"]["Output"] = api_response.text
    if api_response.status_code != 200:
        temp["s"] = "f"
        temp["Reason"] = api_response.text
    else:
        temp["s"] = "s"
    list_data.append(temp)

# 序列化为标准JSON字符串输出,处理中文和None值
dbutils.notebook.exit(json.dumps(list_data, ensure_ascii=False, default=str))

ADF端配置

在ForEach活动的Items字段,使用以下表达式将返回的JSON字符串解析为数组:

@json(activity('test_notebook').output.runOutput)

内容的提问来源于stack exchange,提问作者Developer Rajinikanth

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 13:23:18