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

如何在Databricks Workflow中实现笔记本的循环调用

Databricks Workflow实现循环调用Notebook的思路

以下是几种可行的实现方案,可根据你的业务需求选择:

方案一:动态任务生成(Databricks 1.13+ 推荐)

利用Databricks Workflow的动态任务组功能,通过一个父Notebook生成子任务列表,自动创建多个调用2_CIRCUITO的任务。

步骤:

  1. 创建一个父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}))
  1. 在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 23:53:19