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

AzureML中基于PythonScriptStep的流水线步骤间无落地数据传递方法问询

如何在Azure ML Pipeline中实现步骤间无落地的数据传递

嘿,这个问题我刚好研究过!Azure ML Pipeline默认确实会依赖Blob存储来传递步骤间的数据,但要实现纯内存流转、不把数据存到任何持久化存储的话,有几个实用方案,你可以根据数据体量和步骤独立性需求来选:

1. 最可靠:合并两个步骤为单个脚本

如果你的数据预处理和训练逻辑资源需求一致(比如不需要不同的计算资源),最直接的办法就是把data_prep.py和train.py的逻辑合并到同一个脚本里。这样数据完全在内存中流转,根本不需要跨步骤传递,自然也就不会落地到Blob。

举个合并后的脚本示例:

# combined_script.py
import pandas as pd
from sklearn.model_selection import train_test_split
from sklearn.ensemble import RandomForestClassifier

# 数据预处理逻辑
def prepare_data():
    # 这里假设你从输入数据集读取原始数据
    raw_df = pd.read_csv("raw_data.csv")
    # 你的转换操作,比如清洗、特征工程
    processed_df = raw_df.dropna().drop(columns=["irrelevant_col"])
    return processed_df

# 模型训练逻辑
def train_model(processed_data):
    X = processed_data.drop("target_column", axis=1)
    y = processed_data["target_column"]
    X_train, X_test, y_train, y_test = train_test_split(X, y, test_size=0.2)
    
    model = RandomForestClassifier(n_estimators=100)
    model.fit(X_train, y_train)
    # 这里可以直接在内存中评估模型,或者按需保存模型(模型存储和数据传递是两回事)
    return model

if __name__ == "__main__":
    # 内存中直接传递数据
    cleaned_data = prepare_data()
    trained_model = train_model(cleaned_data)
    print("训练完成!")

然后用单个PythonScriptStep来运行这个脚本即可,完全避免了中间数据的持久化。

2. 小体量数据:通过Pipeline参数传递字符串化数据

如果你的预处理后数据是极小的(比如几个统计参数、小型特征字典),可以把数据序列化为JSON字符串,通过PipelineParameter在步骤间传递。这种方式完全在内存中流转,不会落地到任何存储。

步骤配置示例:

from azureml.pipeline.steps import PythonScriptStep
from azureml.pipeline.core import PipelineParameter
from azureml.core import Workspace, ComputeTarget

ws = Workspace.from_config()
compute_target = ComputeTarget(ws, "your-compute-cluster")

# 定义传递数据的参数
processed_data_param = PipelineParameter(name="processed_data", default_value="")

# 预处理步骤:将数据序列化为字符串并传递
prep_step = PythonScriptStep(
    script_name="data_prep.py",
    arguments=["--output-data", processed_data_param],
    compute_target=compute_target,
    source_directory="./scripts"
)

# 训练步骤:接收字符串并反序列化为数据
train_step = PythonScriptStep(
    script_name="train.py",
    arguments=["--input-data", processed_data_param],
    compute_target=compute_target,
    source_directory="./scripts",
    dependencies=[prep_step]
)

data_prep.py 核心代码:

import pandas as pd
import json
import argparse

parser = argparse.ArgumentParser()
parser.add_argument("--output-data", type=str)
args = parser.parse_args()

# 预处理逻辑
raw_df = pd.read_csv("raw_data.csv")
processed_df = raw_df.dropna()

# 序列化为JSON字符串
data_str = processed_df.to_json(orient="split")
# 这里通过打印输出传递给下一个步骤(实际Azure ML会捕获这个参数传递)
print(f"--output-data {data_str}")

train.py 核心代码:

import pandas as pd
import json
import argparse

parser = argparse.ArgumentParser()
parser.add_argument("--input-data", type=str)
args = parser.parse_args()

# 反序列化字符串为DataFrame
processed_df = pd.read_json(args.input_data, orient="split")
# 后续训练逻辑...

注意:这个方法只适合KB级别的小数据,如果是大型数据集,序列化和传递的开销会非常大,甚至超出参数长度限制。

3. 大数据可选:依赖同一计算节点的本地临时存储

如果必须分开两个步骤,且处理的是大数据,可以尝试让两个步骤运行在同一个计算节点的本地临时目录(比如/tmp)。Azure ML如果启用了集群复用,可能会把两个步骤调度到同一个节点上,数据只存在节点本地内存/临时磁盘,不会上传到Blob,步骤结束后临时数据会被自动清理。

步骤配置示例:

prep_step = PythonScriptStep(
    script_name="data_prep.py",
    arguments=["--output-dir", "/tmp/processed_data"],
    compute_target=compute_target,
    source_directory="./scripts",
    allow_reuse=True  # 启用集群复用
)

train_step = PythonScriptStep(
    script_name="train.py",
    arguments=["--input-dir", "/tmp/processed_data"],
    compute_target=compute_target,
    source_directory="./scripts",
    allow_reuse=True,
    dependencies=[prep_step]
)

data_prep.py 核心代码:

import pandas as pd
import argparse
import os

parser = argparse.ArgumentParser()
parser.add_argument("--output-dir", type=str)
args = parser.parse_args()

# 创建临时目录
os.makedirs(args.output_dir, exist_ok=True)

# 预处理逻辑
raw_df = pd.read_csv("raw_data.csv")
processed_df = raw_df.dropna()

# 写入本地临时目录
processed_df.to_csv(os.path.join(args.output_dir, "processed.csv"), index=False)

train.py 核心代码:

import pandas as pd
import argparse
import os

parser = argparse.ArgumentParser()
parser.add_argument("--input-dir", type=str)
args = parser.parse_args()

# 从本地临时目录读取数据
processed_df = pd.read_csv(os.path.join(args.input_dir, "processed.csv"))
# 后续训练逻辑...

⚠️ 注意:这个方法不保证100%可靠,因为Azure ML的调度器可能会把两个步骤分配到不同的节点,这时候训练步骤就会找不到数据。适合测试场景或者对可靠性要求不高的流水线。

内容的提问来源于stack exchange,提问作者RaLa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 07:07:50