如何基于DataBricks Notebook JSON输出使用ForEach活动并集成Lookup
集成方案与实现步骤
一、修正并预处理Databricks Notebook输出JSON
你提供的示例JSON存在语法错误("vlue2": "value": { 格式非法),且ForEach活动仅支持遍历数组结构,建议先调整输出格式:
方案1:修改Notebook直接输出数组
在Databricks Notebook中,将需要循环的对象整理为数组输出,示例结构如下:
{ "runPageUrl": "https://adb-url", "runOutput": [ { "xxxx": 60, "mmmm": "value", "db": "dbname", "table": "tablename", "asset_id": "1cda102e-dddddxxxxx9c775", "aggvalue": ["value"], "date_val": ["date_val"] }, { "xxxx": 60, "mmmm": "value", "db": "dbname", "table": "tablename", "asset_id": "1cda102e-dddddxxxxx9c775", "aggvalue": ["value"], "date_val": ["date_val"] } ], "effectiveIntegrationRuntime": "AutoResolveIntegrationRuntime", "executionDuration": 297, "durationInQueue": { "integrationRuntimeQueue": 0 }, "billingReference": { "activityType": "ExternalActivity", "billableDuration": [ { "meterType": "AzureIR", "duration": 0.3343434, "unit": "Hours" } ] } }
方案2:ADF内转换现有输出为数组
若无法修改Notebook,添加变量活动,用表达式将runOutput下的零散对象转为数组:
@createArray(activity('DatabricksNotebook1').output.runOutput.value, activity('DatabricksNotebook1').output.runOutput.vlue2)
二、配置ForEach活动遍历输出
- 在ADF管道中添加ForEach活动,设置Items属性为:
(若用变量转换后的数组,替换为@activity('DatabricksNotebook1').output.runOutput@variables('转换后的数组变量名')) - 根据需求选择执行模式:保持默认并行(提高效率)或勾选Sequential(按顺序执行)。
三、ForEach内部调用Lookup活动获取SQL数据
- 在ForEach活动中嵌入Lookup活动,连接目标Azure SQL DB。
- 编写带参数的查询语句,使用ForEach当前循环项的字段作为查询条件,示例:
SELECT detail_col1, detail_col2 FROM your_target_table WHERE asset_id = '@{item().asset_id}' - 若只需要单条匹配数据,将Lookup的First row only属性设为
True。
四、拼接参数并调用第二个Databricks Notebook
- 在Lookup活动后添加Databricks Notebook活动。
- 配置Notebook的Base parameters,将循环项字段与Lookup结果拼接传入,示例参数:
database:@{item().db}table:@{item().table}asset_full_info:@{concat(item().asset_id, '_', activity('LookupSQL').output.firstRow.detail_col1)}agg_list:@{item().aggvalue}
- 若需传入复杂结构,可拼接为JSON字符串后传入:
@string(json(concat('{"db":"', item().db, '","detail":"', activity('LookupSQL').output.firstRow.detail_col1, '"}')))
关键优化与注意事项
- 性能:数据量较大时,建议在Databricks Notebook中直接批量拉取SQL数据,减少ADF Lookup的调用次数;并行执行ForEach时需注意Azure SQL DB的并发连接限制。
- 错误处理:为所有活动添加重试策略,同时用If Condition活动捕获Lookup无结果的场景,避免管道中断。
- 调试:添加Set Variable活动临时存储中间结果,或在Notebook中输出日志,验证参数是否正确传递。
内容的提问来源于stack exchange,提问作者Developer Rajinikanth
相关产品推荐
相关产品推荐

