如何将Pipeline数据传递给Azure ML管线的Databricks步骤
Azure ML 流水线Python步骤与Databricks步骤互传中间数据解决方案
问题根因
Azure ML 给不同计算环境返回的PipelineData路径格式会自动适配当前运行环境:
- Python脚本步骤运行在Azure ML计算实例/集群上,返回本地挂载的POSIX格式路径
- Databricks步骤运行在Databricks集群上,返回
wasbs://协议的云存储路径
两种路径格式无法直接跨环境识别,导致数据读取失败。
具体实现方案
- 替换中间输出定义方式:弃用原生
PipelineData,改用OutputFileDatasetConfig定义中间数据,该类会自动给不同计算环境返回适配的路径格式,无需手动转换:
from azureml.data import OutputFileDatasetConfig # 定义中间输出,路径自动绑定流水线运行ID,避免不同运行的数据冲突 prepped_data_output = OutputFileDatasetConfig( name="prepped_parameter", destination=(data_store, "pipeline_intermediates/{run-id}/prepped_parameter") ).as_upload()
将该对象作为前序Python步骤的输出、后续所有Python/Databricks步骤的输入即可。
- 提前挂载Azure ML存储到Databricks DBFS:在Databricks集群的初始化脚本中添加存储挂载逻辑,一次配置永久生效,无需每次运行流水线手动操作:
// Databricks集群初始化脚本示例 val storageAccountName = "你的Azure存储账户名" val containerName = "Azure ML blob存储容器名" val storageAccountKey = "你的存储账户访问密钥" val configs = Map( s"fs.azure.account.key.$storageAccountName.blob.core.windows.net" -> storageAccountKey ) dbutils.fs.mount( source = s"wasbs://$containerName@$storageAccountName.blob.core.windows.net/", mountPoint = "/mnt/azureml_workspace_store", extraConfigs = configs )
挂载完成后,在Databricks步骤中可以将输入路径直接转换为DBFS挂载路径读取,示例:wasbs://xxx@xxx.blob.core.windows.net/azureml/xxx/ 转换为 /dbfs/mnt/azureml_workspace_store/azureml/xxx/,支持pandas、pyspark等所有常用读取方式。
自动清理中间数据配置:给存储容器配置生命周期管理规则,指定超过指定天数(比如7天)的中间数据自动删除,完全无需手动操作。如果是临时测试流水线,也可以在流水线运行完成后通过Azure ML SDK批量删除对应运行ID路径下的所有文件,仅需一行代码即可完成清理。
Databricks步骤间传数据优化:两个Databricks步骤之间传数据时,直接使用DBFS挂载路径作为参数传递即可,无需走
wasbs协议,读写速度更快,不存在路径识别问题。
内容的提问来源于stack exchange,提问作者Abhijit Halder
相关产品推荐
相关产品推荐

