如何在Databricks Workflow中实现笔记本的循环调用
Databricks Workflow实现循环调用Notebook的思路
以下是几种可行的实现方案,可根据你的业务需求选择:
方案一:动态任务生成(Databricks 1.13+ 推荐)
利用Databricks Workflow的动态任务组功能,通过一个父Notebook生成子任务列表,自动创建多个调用2_CIRCUITO的任务。
步骤:
- 创建一个父Notebook(比如命名为
0_GEN_TASKS),负责生成待执行的任务列表:
# 读取你的id列表(替换成实际数据源,比如从Delta表/数据库读取) list_RREGA = [{"id_in": "id_001"}, {"id_in": "id_002"}, {"id_in": "id_003"}] # 构建动态任务配置 tasks = [] for idx, doc in enumerate(list_RREGA): task_key = f"process_{doc['id_in']}_{idx}" tasks.append({ "task_key": task_key, "notebook_task": { "notebook_path": "/Workspace/你的路径/2_CIRCUITO", "base_parameters": {"id_in": doc["id_in"]} }, # 可复用集群或指定新集群,根据需求调整 "existing_cluster_id": "你的集群ID" }) # 输出JSON格式的任务列表,供Workflow解析 import json print(json.dumps({"tasks": tasks}))
- 在Databricks Workflow中创建任务组,选择「动态任务」类型,将生成器任务设置为上述父Notebook。Workflow会自动解析父Notebook输出的JSON,生成对应数量的子任务,每个子任务独立调用
2_CIRCUITO并传入参数。
方案二:通过API/CLI批量创建任务
如果需要自动化部署或更灵活的控制,可使用Databricks API或CLI直接生成多个调用2_CIRCUITO的任务,组成一个作业。
示例(Python调用API):
import requests import json # 配置Databricks API信息 databricks_host = "你的Databricks工作区URL" token = "你的个人访问令牌" headers = {"Authorization": f"Bearer {token}"} # 读取id列表 list_RREGA = [{"id_in": "id_001"}, {"id_in": "id_002"}] # 构建作业任务列表 tasks = [] for doc in list_RREGA: tasks.append({ "task_key": f"run_{doc['id_in']}", "notebook_task": { "notebook_path": "/Workspace/你的路径/2_CIRCUITO", "base_parameters": {"id_in": doc["id_in"]} }, "existing_cluster_id": "你的集群ID" }) # 提交作业请求 job_payload = { "name": "批量处理CIRCUITO任务", "tasks": tasks, "email_notifications": {"on_failure": ["你的邮箱"]} } response = requests.post(f"{databricks_host}/api/2.1/jobs/create", headers=headers, json=job_payload) print(response.json())
方案三:单任务内并行/串行执行(简化版)
如果不需要单独监控每个子任务,可保留原有的循环逻辑,将1_CIRCUITO作为Workflow的一个独立任务,在任务内部优化执行效率(比如用线程池并行调用)。
优化后的代码示例:
from concurrent.futures import ThreadPoolExecutor list_RREGA = [{"id_in": "id_001"}, {"id_in": "id_002"}] # 定义调用Notebook的函数 def run_notebook(doc): dbutils.notebook.run("2_CIRCUITO", 600, {"id_in": doc["id_in"]}) # 并行执行(根据集群资源调整线程数) with ThreadPoolExecutor(max_workers=5) as executor: executor.map(run_notebook, list_RREGA)
将这个Notebook作为Workflow的一个任务,利用Workflow的调度、失败重试等功能,同时通过并行执行提升效率。
内容的提问来源于stack exchange,提问作者Jose Maria Acevedo Reinoso
相关产品推荐
相关产品推荐

