Azure ML SDK V1迁移V2:无输入输出Pipeline编排求助
Azure ML SDK V1 Pipeline迁移至V2:无输入输出Pipeline编排问题
问题概述
正在将Azure Machine Learning SDK V1中的Pipeline迁移至V2,因对V2逻辑理解不足受阻。V1中仅需创建PythonScriptStep并通过StepSequence编排即可部署无输入输出的Pipeline(数据存储于ADLS Gen2,使用Databricks表作为输入),但V2中无法直接实现两步连续执行,尝试自定义CommandSequence时出现报错。
V1实现代码
script_step_1 = PythonScriptStep( name="step1", script_name="main.py", arguments=arguments, # list of PipelineParameter compute_target=ComputeTarget(workspace=ws, name="cpu-16-128"), source_directory="./my_project_folder", runconfig=runconfig, # Conda + extra index url + custom dockerfile allow_reuse=False, ) script_step_2 = PythonScriptStep( name="step2", ... ) step_sequence = StepSequence( steps=[ script_step_1, script_step_2, ] ) # Create Pipeline pipeline = Pipeline( workspace=ws, steps=step_sequence, ) pipeline_run = experiment.submit(pipeline)
V2当前尝试
已通过BuildContext结合Dockerfile创建环境,并定义了调用Python脚本的command组件:
环境定义
azureml_env = Environment( build=BuildContext( path="./docker_folder", # With Dockerfile and requirements.txt ), name="my-project-env", )
Command组件定义
step_1 = command( environment=azureml_env , command="python main.py", code="./my_project_folder", ) # step_2定义类似step_1
错误的Pipeline编排尝试
@pipeline(compute="serverless") def default_pipeline(): return { "my_pipeline": [step_1, step_2] } # 提交Pipeline my_pipeline = default_pipeline() pipeline_job = ml_client.jobs.create_or_update( my_pipeline, experiment_name=experiment_name, )
自定义CommandSequence的报错问题
尝试通过添加虚拟输入输出实现编排,但报错AttributeError: 'dict' object has no attribute 'my_output':
自定义CommandSequence代码
class CommandSequence: def __init__(self, commands, ml_client): self.commands = commands self.ml_client = ml_client def build(self): for i in range(len(self.commands)): cmd = self.commands[i] if i == 0: cmd = command( display_name=cmd.display_name, description=cmd.description, environment=cmd.environment, command=cmd.command, code=cmd.code, is_deterministic=cmd.is_deterministic, outputs=dict( my_output=Output(type="uri_folder", mode="rw_mount"), ), ) else: cmd = command( display_name=cmd.display_name, description=cmd.description, environment=cmd.environment, command=cmd.command, code=cmd.code, is_deterministic=cmd.is_deterministic, inputs=self.commands[i - 1].outputs.my_output, outputs=dict( my_output=Output(type="uri_folder", mode="rw_mount"), ), ) cmd = self.ml_client.create_or_update(cmd.component) self.commands[i] = cmd print(self.commands[i]) return self.commands
调用代码
@pipeline(compute="serverless") def default_pipeline(): command_sequence = CommandSequence([step_1, step_2], ml_client).build() return { "my_pipeline": command_sequence[-1].outputs.my_output }
解决方案
方案1:无需自定义类,直接实现串行执行
V2中无需依赖虚拟输入输出,可通过显式添加依赖关系让step2在step1完成后执行:
@pipeline(compute="serverless") def default_pipeline(): # 执行第一步 step1_exec = step_1() # 执行第二步 step2_exec = step_2() # 添加依赖:step2必须在step1完成后运行 step2_exec.add_dependency(step1_exec) # 无输出需求时返回空字典即可 return {}
若偏好通过输入输出绑定实现依赖(更符合V2的组件化逻辑),可给组件添加虚拟输入输出:
# 修改step_1,添加虚拟输出 step_1 = command( environment=azureml_env, command="python main.py", code="./my_project_folder", outputs=dict(dummy_output=Output(type="uri_folder")) ) # 修改step_2,添加虚拟输入 step_2 = command( environment=azureml_env, command="python main.py", code="./my_project_folder", inputs=dict(dummy_input=Input(type="uri_folder")) ) # 编排Pipeline时绑定依赖 @pipeline(compute="serverless") def default_pipeline(): s1 = step_1() # 通过绑定虚拟输入输出,自动形成串行依赖 s2 = step_2(dummy_input=s1.outputs.dummy_output) return {}
方案2:修复自定义CommandSequence的报错
报错原因是:组件的outputs是字典类型,需通过键索引而非属性访问;同时command的inputs参数需要传入字典格式的键值对。修改后的build方法如下:
def build(self): for i in range(len(self.commands)): cmd = self.commands[i] if i == 0: cmd = command( display_name=cmd.display_name, description=cmd.description, environment=cmd.environment, command=cmd.command, code=cmd.code, is_deterministic=cmd.is_deterministic, outputs=dict( my_output=Output(type="uri_folder", mode="rw_mount"), ), ) else: # 1. 用键索引访问outputs,而非属性 # 2. inputs传入字典格式的键值对 cmd = command( display_name=cmd.display_name, description=cmd.description, environment=cmd.environment, command=cmd.command, code=cmd.code, is_deterministic=cmd.is_deterministic, inputs=dict(dummy_input=self.commands[i - 1].outputs["my_output"]), outputs=dict( my_output=Output(type="uri_folder", mode="rw_mount"), ), ) cmd = self.ml_client.create_or_update(cmd.component) self.commands[i] = cmd return self.commands
修改后,调用CommandSequence即可正确绑定step间的依赖关系。
内容的提问来源于stack exchange,提问作者BeGreen
相关产品推荐
相关产品推荐

