Azure DataBricks输出无法在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要求的数组输入格式。
修复步骤
- 引入
json模块,把Python列表序列化为标准JSON字符串,确保格式符合ADF解析要求 - 处理Python的
None值,自动转换为JSON的null - 修正代码里的拼写错误(比如
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

