ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

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

2026/9/18 10:00:25 拓冰建站 浏览量
ZenML Databricks Orchestrator 实战指南:在 Databricks 上编排你的 ML 流水线 ZenML Databricks Orchestrator 实战指南在 Databricks 上编排你的 ML 流水线【免费下载链接】zenmlZenML : One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenmlDatabricks 是统一的数据分析平台将数据仓库与数据湖的优势结合为大数据处理与机器学习提供一体化解决方案。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-sdk0.28.0并随 numpy、pandas 依赖一起安装。同一 integration 还同时提供databricks的 step operator 与 model deployer flavor因此一个 integration 可以支撑编排 部分步骤 模型部署的组合场景。前置条件开始使用前你需要准备一个可用的 Databricks workspace注意文档中链接的 Databricks 官方云平台开通指引分别针对 AWS、Azure、GCP请按你的云环境选择AWSDatabricks 账户开通指南Azure创建 Azure Databricks workspace 指南GCPGCP Databricks 入门指南。一个有权限创建和运行 job 的 Databricks 账户或 service account。官方建议创建一个专用的 Databricks service account为其生成client_id和client_secret用于 API 认证。另外从 orchestrator 的 stack validator 可以看到一条重要的架构约束Databricks orchestrator 在远端运行流水线因此 stack 中所有组件都必须是远程组件。如果某个组件如 artifact store、container registry被标记为 localstack validator 会直接拒绝提交报错提示该组件在 Databricks step 中不可用。工作原理当你使用 Databricks orchestrator 运行流水线时ZenML 会执行以下步骤构建 Python wheelZenML 将你的项目代码打包成一个 Python wheel通过WheeledOrchestrator.create_wheel见 databricks_orchestrator.py 中的submit_pipeline。上传 wheel 到 Databricks workspacewheel 被上传到/Workspace/Shared/.zenml/pipeline/run_namespace/orchestrator/目录前缀常量DATABRICKS_WHEELS_DIRECTORY_PREFIX /Workspace/Shared/.zenml定义在 databricks_utils.py。若上传失败或 job 提交失败ZenML 会尝试递归删除该目录做清理delete_workspace_directory。创建 Databricks jobZenML 通过 Databricks SDKWorkspaceClient创建一个 jobjob 中的每个 task 对应流水线中的一个 steptask 间的depends_on关系精确镜像 step 的 upstream 依赖见_construct_databricks_pipeline。运行 jobjob 创建成功后立即通过jobs.run_now(job_id...)触发运行。job 使用的集群配置完全来自 orchestrator settings包括 Spark 版本、worker 数量或自动伸缩、节点类型以及任意 Spark 配置。Databricks 启动 job 时每个 task 会安装上传的 wheel 并执行对应的 ZenML step entrypointPythonWheelTask(package_namezenml, entry_pointentrypoint.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_package与databricks_job_id作为短参数传入再在运行时重构长环境变量并将ZENML_DATABRICKS_ORCHESTRATOR_RUN_ID写入环境供get_orchestrator_run_id读取。如何使用1. 安装 Integration首先安装 Databricks integrationzenml integration install databricks该命令会安装databricks-sdk以及配套的 numpy、pandas 依赖。2. 注册 Orchestrator 并配置认证注册 orchestrator 时需要指定--flavordatabricks、workspace 的--host以及 service principal 的client_id/client_secret。推荐使用 ZenML secret 引用{{secret.key}}语法避免明文暴露凭据zenml orchestrator register databricks_orchestrator \ --flavordatabricks \ --hosthttps://xxxxx.x.azuredatabricks.net \ --client_id{{databricks.client_id}} \ --client_secret{{databricks.client_secret}}认证细节推荐创建具备创建与运行 job 权限的 Databricks service account然后为其生成client_id和client_secret进行认证Databricks 官方文档提供 service account 创建方式以及如上的权限配置截图。源码中的DatabricksOrchestratorConfig见 databricks_orchestrator_flavor.py会校验client_id与client_secret必须同时提供或同时不提供二者只配置其一会在模型校验阶段直接抛错。若两者都未配置_get_databricks_client会退回仅凭host构造客户端此时依赖 Databricks 环境中的其他认证方式如环境变量、profile 或 service connector。此外DatabricksOrchestratorConfig.is_remote恒为True、is_schedulable恒为True即这是一个支持调度的远程 orchestrator。3. 注册 Stack 并运行将 orchestrator 加入 stack省略号处补充 artifact store 等其他远程组件并设为 activezenml stack register databricks_stack -o databricks_orchestrator ... --set然后像往常一样运行流水线python run.pyDatabricks UIDatabricks 自带完整的 UI你可以用它查看 pipeline run 的更多细节例如每个 step 的日志对于任何在 Databricks 上执行的 run你可以通过以下 Python 代码获取指向 Databricks UI 的 URLfrom 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( scheduleSchedule( cron_expression*/5 * * * * ) )⚠️ Databricks orchestrator只支持Schedule对象中的cron_expression字段传入的其他调度参数如interval_second、catchup等都会被忽略——源码会在检测到这些参数时输出 warning 日志见submit_pipeline中的校验逻辑。⚠️ 使用 cron 调度时必须通过 orchestrator settings 中的schedule_timezone指定一个合法的 IANA 时区 ID例如America/New_York或UTC。如果配置了cron_expression却未设置schedule_timezone提交会在两个层面被拦截settings 校验器要求ZoneInfo(value)能正确解析时区见DatabricksOrchestratorSettings._validate_schedule_timezonesubmit_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.12num_workersint 0固定 worker 数量不能与autoscale同时使用node_type_idstr默认Standard_D4s_v5Databricks 节点类型标识参考官方实例类型文档driver_node_type_idstrSpark driver 的节点类型不指定时默认与 worker 相同policy_idstrDatabricks cluster policy ID用于治理与成本控制。未指定时会尝试查找名为Job Compute的默认 policyautoscale(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 attributescustom_tagsDict[str, str]最多 45 个应用到底层集群资源如 AWS EC2 实例、EBS 卷的标签用于成本分摊与治理job_tagsDict[str, str]最多 25 个应用到 Databricks job 本身并转发为集群标签access_control_listList[DatabricksAccessControlRequest]job 的访问控制列表可授予用户/组/service principal 权限CAN_VIEW、CAN_MANAGE_RUN、CAN_MANAGE、IS_OWNER。默认只有 job 创建者可访问。每条 ACL 必须恰好指定一个主体timeout_secondsint 00 表示无超时job 每次 run 的超时时间task_timeout_secondsint 00 表示无超时job 中每个 taskstep的超时时间init_scriptsList[str]集群初始化脚本只支持以dbfs:/开头的 DBFS 路径docker_image_urlstr集群使用的 Docker 镜像 URL需可从 Databricks workspace 访问docker_image_username/docker_image_passwordstrDocker registry 认证凭据必须成对提供否则校验报错schedule_timezonestrIANA 时区定时执行使用的时区仅在配置 cron 调度时使用max_concurrent_runsint11000job 的最大并发 run 数Databricks 未指定时默认为 1max_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_version15.3.x-scala2.12, num_workers3, node_type_idStandard_D4s_v5, policy_idPOLICY_ID, spark_conf{}, spark_env_vars{}, init_scripts[dbfs:/scripts/install_dependencies.sh], schedule_timezoneAmerica/Los_Angeles, )固定大小集群 vs 自动伸缩集群使用num_workers表示固定大小集群需要自动伸缩时省略num_workers并设置autoscale例如autoscale(2, 3)。num_workers与autoscale二选一的规则在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_version15.3.x-scala2.12, num_workers3, node_type_idStandard_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_version和node_type_id设置为支持 GPU 的值即可使用 GPU 集群from zenml.integrations.databricks.flavors.databricks_orchestrator_flavor import DatabricksOrchestratorSettings databricks_settings DatabricksOrchestratorSettings( spark_version15.3.x-gpu-ml-scala2.12, node_type_idStandard_NC24ads_A100_v4, policy_idPOLICY_ID, autoscale(1, 2), )通过这些设置orchestrator 会使用 GPU 版本的 Spark 和 GPU 节点类型。为 GPU 硬件启用 CUDA如果 step 中需要 CUDA请参考 ZenML 的分布式训练指南来配置所需的依赖与运行时设置。源码级验证以上行为均有测试佐证集成测试 test_databricks_orchestrator.py 使用 mock 的 Databricks client验证了 orchestrator 会把access_control_list、availability_type、custom_tags、docker_image_url与凭据、init_scripts、job_tags、max_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),仅供参考