ZenML Databricks Orchestrator 实战指南:在 Databricks 上编排你的 ML 流水线
2026/9/18 10:00:20 网站建设 项目流程

ZenML Databricks Orchestrator 实战指南:在 Databricks 上编排你的 ML 流水线

【免费下载链接】zenmlZenML 🙏: One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml

Databricks 是统一的数据分析平台,将数据仓库与数据湖的优势结合,为大数据处理与机器学习提供一体化解决方案。ZenML 的 Databricks integration 提供了一种名为databricks的 orchestrator flavor,让你能够在 ZenML 框架内把完整流水线直接提交到 Databricks 上运行,借助其分布式计算能力和针对大数据、机器学习优化过的运行环境。读完本文,你将掌握 Databricks orchestrator 的注册配置、认证方式、集群与任务参数调优、定时调度,以及如何在 Databricks UI 中追踪运行状态。

⚠️Alpha 功能提醒:本文涉及的 Databricks orchestrator 部分能力目前处于 Alpha 阶段,后续可能发生变化。建议在受控环境中使用,并向 ZenML 团队反馈问题。

💡适用范围:如果只想把部分 step 放到 Databricks 上执行、整体流水线仍由其他 orchestrator 编排,请改用 Databricks step operator,而不是本 orchestrator。

何时使用 Databricks Orchestrator

在以下场景中,你应该考虑使用 Databricks orchestrator:

  • 你已经在使用 Databricks 承载数据和 ML 工作负载;
  • 你希望利用 Databricks 强大的分布式计算能力运行 ML 流水线;
  • 你在寻找一个与 Databricks 其他服务(SQL、Delta Lake、MLflow 等)集成良好的托管方案;
  • 你想借助 Databricks 针对大数据处理和机器学习的优化能力。

从源码层面看,该 orchestrator 由DatabricksIntegration注册(见 src/zenml/integrations/databricks/init.py),依赖固定为databricks-sdk==0.28.0,并随 numpy、pandas 依赖一起安装。同一 integration 还同时提供databricks的 step operator 与 model deployer flavor,因此一个 integration 可以支撑"编排 + 部分步骤 + 模型部署"的组合场景。

前置条件

开始使用前,你需要准备:

  1. 一个可用的 Databricks workspace(注意:文档中链接的 Databricks 官方云平台开通指引分别针对 AWS、Azure、GCP,请按你的云环境选择):
    • AWS:Databricks 账户开通指南;
    • Azure:创建 Azure Databricks workspace 指南;
    • GCP:GCP Databricks 入门指南。
  2. 一个有权限创建和运行 job 的 Databricks 账户或 service account。官方建议创建一个专用的 Databricks service account,为其生成client_idclient_secret用于 API 认证。

另外,从 orchestrator 的 stack validator 可以看到一条重要的架构约束:Databricks orchestrator 在远端运行流水线,因此 stack 中所有组件都必须是远程组件。如果某个组件(如 artifact store、container registry)被标记为 local,stack validator 会直接拒绝提交,报错提示该组件"在 Databricks step 中不可用"。

工作原理

当你使用 Databricks orchestrator 运行流水线时,ZenML 会执行以下步骤:

  1. 构建 Python wheel:ZenML 将你的项目代码打包成一个 Python wheel(通过WheeledOrchestrator.create_wheel,见 databricks_orchestrator.py 中的submit_pipeline)。
  2. 上传 wheel 到 Databricks workspace:wheel 被上传到/Workspace/Shared/.zenml/<pipeline>/<run_namespace>/orchestrator/目录(前缀常量DATABRICKS_WHEELS_DIRECTORY_PREFIX = "/Workspace/Shared/.zenml"定义在 databricks_utils.py)。若上传失败或 job 提交失败,ZenML 会尝试递归删除该目录做清理(delete_workspace_directory)。
  3. 创建 Databricks job:ZenML 通过 Databricks SDK(WorkspaceClient)创建一个 job,job 中的每个 task 对应流水线中的一个 step,task 间的depends_on关系精确镜像 step 的 upstream 依赖(见_construct_databricks_pipeline)。
  4. 运行 job:job 创建成功后立即通过jobs.run_now(job_id=...)触发运行。

job 使用的集群配置完全来自 orchestrator settings,包括 Spark 版本、worker 数量或自动伸缩、节点类型以及任意 Spark 配置。Databricks 启动 job 时,每个 task 会安装上传的 wheel 并执行对应的 ZenML step entrypoint(PythonWheelTask(package_name="zenml", entry_point="entrypoint.main"),见convert_step_to_task)。

📌关于 wheel 的保留策略:由于 Databricks 的 job 定义、定时调度以及手动重跑都会持续引用 workspace 中的 wheel 文件,orchestrator 在 job 成功提交后不会删除/Workspace/Shared/.zenml下的 wheel 包。请根据团队保留策略定期清理该路径下不再需要重跑的旧 wheel 目录,避免 workspace 空间膨胀。

在任务运行阶段,每个 task 通过DatabricksEntrypointConfiguration(见 databricks_orchestrator_entrypoint_config.py)来恢复运行环境:由于 Databricks job 参数有 256 字符长度限制,该入口配置会把wheel_packagedatabricks_job_id作为短参数传入,再在运行时重构长环境变量,并将ZENML_DATABRICKS_ORCHESTRATOR_RUN_ID写入环境,供get_orchestrator_run_id读取。

如何使用

1. 安装 Integration

首先安装 Databricks integration:

zenml integration install databricks

该命令会安装databricks-sdk以及配套的 numpy、pandas 依赖。

2. 注册 Orchestrator 并配置认证

注册 orchestrator 时,需要指定--flavor=databricks、workspace 的--host以及 service principal 的client_id/client_secret。推荐使用 ZenML secret 引用({{secret.key}}语法)避免明文暴露凭据:

zenml orchestrator register databricks_orchestrator \ --flavor=databricks \ --host="https://xxxxx.x.azuredatabricks.net" \ --client_id={{databricks.client_id}} \ --client_secret={{databricks.client_secret}}

💡认证细节:推荐创建具备创建与运行 job 权限的 Databricks service account,然后为其生成client_idclient_secret进行认证(Databricks 官方文档提供 service account 创建方式,以及如上的权限配置截图)。

源码中的DatabricksOrchestratorConfig(见 databricks_orchestrator_flavor.py)会校验:client_idclient_secret必须同时提供或同时不提供,二者只配置其一会在模型校验阶段直接抛错。若两者都未配置,_get_databricks_client会退回仅凭host构造客户端,此时依赖 Databricks 环境中的其他认证方式(如环境变量、profile 或 service connector)。此外,DatabricksOrchestratorConfig.is_remote恒为Trueis_schedulable恒为True,即这是一个支持调度的远程 orchestrator。

3. 注册 Stack 并运行

将 orchestrator 加入 stack(省略号处补充 artifact store 等其他远程组件),并设为 active:

zenml stack register databricks_stack -o databricks_orchestrator ... --set

然后像往常一样运行流水线:

python run.py

Databricks UI

Databricks 自带完整的 UI,你可以用它查看 pipeline run 的更多细节,例如每个 step 的日志:

对于任何在 Databricks 上执行的 run,你可以通过以下 Python 代码获取指向 Databricks UI 的 URL:

from zenml.client import Client pipeline_run = Client().get_pipeline_run("<PIPELINE_RUN_NAME>") orchestrator_url = pipeline_run.run_metadata["orchestrator_url"].value

该 URL 由 orchestrator 在get_pipeline_run_metadata中生成,格式为{host}/jobs/{orchestrator_run_id}orchestrator_run_id即运行时的{{job.id}}参数,见 databricks_orchestrator.py),并以Uri类型的 metadata 写入 run,因此可以直接在 ZenML 中关联到 Databricks 的 job 页面。

定时调度流水线

Databricks orchestrator 支持使用 Databricks 原生的调度能力(Jobs 调度)按计划运行流水线。

如何创建调度

from zenml.config.schedule import Schedule # 每 5 分钟运行一次流水线 pipeline_instance.run( schedule=Schedule( cron_expression="*/5 * * * *" ) )

⚠️ Databricks orchestrator只支持Schedule对象中的cron_expression字段,传入的其他调度参数(如interval_secondcatchup等)都会被忽略——源码会在检测到这些参数时输出 warning 日志(见submit_pipeline中的校验逻辑)。

⚠️ 使用 cron 调度时,必须通过 orchestrator settings 中的schedule_timezone指定一个合法的 IANA 时区 ID(例如America/New_YorkUTC)。如果配置了cron_expression却未设置schedule_timezone,提交会在两个层面被拦截:settings 校验器要求ZoneInfo(value)能正确解析时区(见DatabricksOrchestratorSettings._validate_schedule_timezone),submit_pipeline_upload_and_run_pipeline也会在缺失时区时抛出ValueError

如何删除调度

ZenML 负责创建 Databricks schedule,但调度的生命周期由你在 Databricks 侧管理。要取消一个已调度的 Databricks 流水线,请在 Databricks UI 或 CLI 中删除对应的 schedule。

进阶配置:DatabricksOrchestratorSettings

对于更精细的控制,可以在 pipeline 或 step 级别传入DatabricksOrchestratorSettings。它继承自DatabricksBaseSettings(见 databricks_shared_settings.py),并额外提供调度与 job 级参数:

参数类型 / 默认值说明
spark_versionstr,默认使用工作区默认值(utils 中默认16.4.x-scala2.12Databricks 集群的 Apache Spark 版本,例如"15.3.x-scala2.12"
num_workersint >= 0固定 worker 数量;不能与autoscale同时使用
node_type_idstr,默认Standard_D4s_v5Databricks 节点类型标识,参考官方实例类型文档
driver_node_type_idstrSpark driver 的节点类型,不指定时默认与 worker 相同
policy_idstrDatabricks cluster policy ID,用于治理与成本控制。未指定时会尝试查找名为Job Compute的默认 policy
autoscale(int, int),默认(0, 1)集群自动伸缩的(min_workers, max_workers)边界
autotermination_minutesint >= 0空闲自动终止分钟数,用于控制闲置集群成本
single_user_namestr单用户集群访问模式下的 Databricks 用户名
spark_confDict[str, str]自定义 Spark 配置键值对,例如{"spark.sql.adaptive.enabled": "true"}
spark_env_varsDict[str, str]Spark driver 与 executor 的环境变量
availability_typeON_DEMAND/SPOT/SPOT_WITH_FALLBACK实例可用性类型:按需(有保障)、竞价(成本优化)、竞价带按需兜底。SDK 会根据 host 自动判断云厂商并映射到对应的 AWS/Azure/GCP attributes
custom_tagsDict[str, str],最多 45 个应用到底层集群资源(如 AWS EC2 实例、EBS 卷)的标签,用于成本分摊与治理
job_tagsDict[str, str],最多 25 个应用到 Databricks job 本身并转发为集群标签
access_control_listList[DatabricksAccessControlRequest]job 的访问控制列表,可授予用户/组/service principal 权限(CAN_VIEWCAN_MANAGE_RUNCAN_MANAGEIS_OWNER)。默认只有 job 创建者可访问。每条 ACL 必须恰好指定一个主体
timeout_secondsint >= 0,0 表示无超时job 每次 run 的超时时间
task_timeout_secondsint >= 0,0 表示无超时job 中每个 task(step)的超时时间
init_scriptsList[str]集群初始化脚本,只支持以dbfs:/开头的 DBFS 路径
docker_image_urlstr集群使用的 Docker 镜像 URL,需可从 Databricks workspace 访问
docker_image_username/docker_image_passwordstrDocker registry 认证凭据,必须成对提供,否则校验报错
schedule_timezonestr(IANA 时区)定时执行使用的时区,仅在配置 cron 调度时使用
max_concurrent_runsint,1~1000job 的最大并发 run 数,Databricks 未指定时默认为 1
max_retriesint >= -1,-1 表示无限重试失败 task 的最大重试次数
min_retry_interval_millisint >= 0重试之间的最小间隔(毫秒),例如60000表示间隔 1 分钟
retry_on_timeoutbooltask 超时时是否重试,需同时设置max_retries

一个综合示例:

from zenml.integrations.databricks.flavors.databricks_orchestrator_flavor import DatabricksOrchestratorSettings databricks_settings = DatabricksOrchestratorSettings( spark_version="15.3.x-scala2.12", num_workers=3, node_type_id="Standard_D4s_v5", policy_id=POLICY_ID, spark_conf={}, spark_env_vars={}, init_scripts=["dbfs:/scripts/install_dependencies.sh"], schedule_timezone="America/Los_Angeles", )

固定大小集群 vs 自动伸缩集群:使用num_workers表示固定大小集群;需要自动伸缩时,省略num_workers并设置autoscale,例如autoscale=(2, 3)num_workersautoscale二选一的规则在build_databricks_cluster_spec中体现:设置了num_workers就构造固定 worker 数,否则构造AutoScale对象。默认的autoscale=(0, 1)是刻意为之——它允许 driver-only 集群(省钱),又能在需要时启动一个 worker。

这些 settings 可以同时指定在 pipeline 级别或 step 级别:

# 在 pipeline 级别指定 @pipeline( settings={ "orchestrator": databricks_settings, } ) def my_pipeline(): ...

给 Databricks 资源打标签

为满足成本分摊、治理与项目追踪需求,可以用两个 settings 给 Databricks 资源打标签:

  • custom_tags:应用到底层集群资源(如 AWS EC2 实例、EBS 卷),最多 45 个标签
  • job_tags:应用到 Databricks job 本身,并会被转发为集群标签,最多 25 个标签
from zenml.integrations.databricks.flavors.databricks_orchestrator_flavor import DatabricksOrchestratorSettings databricks_settings = DatabricksOrchestratorSettings( spark_version="15.3.x-scala2.12", num_workers=3, node_type_id="Standard_D4s_v5", custom_tags={"cost_center": "ml-team", "environment": "production"}, job_tags={"project": "recommendation-engine", "owner": "data-team"}, )

注:Databricks 的 tag 值只允许字母数字、下划线、连字符与句点,且最长 63 个字符。utils 中的sanitize_labels会自动将不合规的字符替换为下划线并裁剪长度(见 databricks_utils.py)。

使用 GPU 集群

spark_versionnode_type_id设置为支持 GPU 的值即可使用 GPU 集群:

from zenml.integrations.databricks.flavors.databricks_orchestrator_flavor import DatabricksOrchestratorSettings databricks_settings = DatabricksOrchestratorSettings( spark_version="15.3.x-gpu-ml-scala2.12", node_type_id="Standard_NC24ads_A100_v4", policy_id=POLICY_ID, autoscale=(1, 2), )

通过这些设置,orchestrator 会使用 GPU 版本的 Spark 和 GPU 节点类型。

为 GPU 硬件启用 CUDA

如果 step 中需要 CUDA,请参考 ZenML 的分布式训练指南来配置所需的依赖与运行时设置。

源码级验证

以上行为均有测试佐证:集成测试 test_databricks_orchestrator.py 使用 mock 的 Databricks client,验证了 orchestrator 会把access_control_listavailability_typecustom_tagsdocker_image_url与凭据、init_scriptsjob_tagsmax_concurrent_runs等 settings 完整转发到 job 与集群 payload 中;settings 校验规则(时区合法性、autoscale 边界、init script 的dbfs:/前缀、Docker 凭据成对、service principal 凭据成对)则由 test_databricks_settings.py 覆盖。如果你希望进一步调优,SDK 文档列出了zenml.integrations.databricks下所有可配置属性,ZenML 的步骤与流水线配置文档则解释了 settings 的完整指定方式。

小结

Databricks orchestrator 让 ZenML 流水线得以整体运行在 Databricks 的托管 Spark 环境中:ZenML 负责把项目打成 wheel、上传到 workspace、按 step 依赖创建并触发 Databricks job,同时把集群的 Spark 版本、节点类型、伸缩策略、初始化脚本、Docker 镜像、标签与调度等细节全部暴露为声明式 settings。无论你是想迁移既有 Databricks 工作负载、利用分布式计算加速训练,还是希望复用 Databricks 生态的托管能力,都可以按照本文的步骤快速接入,并通过orchestrator_url元数据把 ZenML 运行视图与 Databricks UI 打通。

【免费下载链接】zenmlZenML 🙏: One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询