【问题标题】:How to pass data or files between Kubeflow containerized components in python如何在 python 中的 Kubeflow 容器化组件之间传递数据或文件
【发布时间】:2020-01-28 17:08:05
【问题描述】:

我正在探索将 Kubeflow 作为部署和连接典型 ML 管道的各种组件的选项。我正在使用 docker 容器作为 Kubeflow 组件,到目前为止,我一直无法成功使用 ContainerOp.file_outputs 对象在组件之间传递结果。

根据我对该功能的理解,创建并保存到声明为组件的 file_outputs 之一的文件应该会导致它持久存在并可供以下组件读取。

这就是我尝试在我的管道 python 代码中声明它的方式:

import kfp.dsl as dsl 
import kfp.gcp as gcp

@dsl.pipeline(name='kubeflow demo')
def pipeline(project_id='kubeflow-demo-254012'):
    data_collector = dsl.ContainerOp(
        name='data collector', 
        image='eu.gcr.io/kubeflow-demo-254012/data-collector',
        arguments=[ "--project_id", project_id ],
        file_outputs={ "output": '/output.txt' }
    )   
    data_preprocessor = dsl.ContainerOp(
        name='data preprocessor',
        image='eu.gcr.io/kubeflow-demo-254012/data-preprocessor',
        arguments=[ "--project_id", project_id ]
    )
    data_preprocessor.after(data_collector)
    #TODO: add other components
if __name__ == '__main__':
    import kfp.compiler as compiler
    compiler.Compiler().compile(pipeline, __file__ + '.tar.gz')

data-collector.py 组件的python 代码中,我获取数据集,然后将其写入output.txt。我可以从同一组件内的文件中读取,但不能在 data-preprocessor.py 内读取,而我得到 FileNotFoundError

file_outputs 的使用对于基于容器的 Kubeflow 组件是否无效,还是我在代码中错误地使用了它?如果在我的情况下不是一个选项,是否可以在管道声明 python 代码中以编程方式创建 Kubernetes 卷并使用它们而不是 file_outputs

【问题讨论】:

    标签: kubeflow kubeflow-pipelines


    【解决方案1】:

    在一个 Kubeflow 管道组件中创建的文件是容器本地的。要在后续步骤中引用它,您需要将其传递为:

    data_preprocessor = dsl.ContainerOp(
            name='data preprocessor',
            image='eu.gcr.io/kubeflow-demo-254012/data-preprocessor',
            arguments=["--fetched_dataset", data_collector.outputs['output'],
                       "--project_id", project_id,
                      ]
    

    注意:data_collector.outputs['output'] 将包含文件/output.txt 的实际字符串内容(不是文件的路径)。如果您希望它包含文件的路径,则需要将数据集写入共享存储(如 s3 或挂载的 PVC 卷)并将共享存储的路径/链接写入/output.txt。然后data_preprocessor 可以根据路径读取数据集。

    【讨论】:

    【解决方案2】:

    主要分为三个步骤:

    1. 保存一个 output.txt 文件,其中将包含您想要传递给下一个组件的数据/参数/任何内容。 注意:它应该在根级别,即/output.txt

    2. 将 file_outputs={'output': '/output.txt'} 作为参数传递,如示例所示。

    3. 在 container_op 中,您将在 dsl.pipeline 中写入传递参数(到需要从早期组件输出的组件的相应参数)作为 comp1.output(这里 comp1 是第一个组件,它产生输出并将其存储在 /输出.txt)

    import kfp
    from kfp import dsl
    
    def SendMsg(
        send_msg: str = 'akash'
    ):
        return dsl.ContainerOp(
            name = 'Print msg', 
            image = 'docker.io/akashdesarda/comp1:latest', 
            command = ['python', 'msg.py'],
            arguments=[
                '--msg', send_msg
            ],
            file_outputs={
                'output': '/output.txt',
            }
        )
    
    def GetMsg(
        get_msg: str
    ):
        return dsl.ContainerOp(
            name = 'Read msg from 1st component',
            image = 'docker.io/akashdesarda/comp2:latest',
            command = ['python', 'msg.py'],
            arguments=[
                '--msg', get_msg
            ]
        )
    
    @dsl.pipeline(
        name = 'Pass parameter',
        description = 'Passing para')
    def  passing_parameter(send_msg):
        comp1 = SendMsg(send_msg)
        comp2 = GetMsg(comp1.output)
    
    
    if __name__ == '__main__':
      import kfp.compiler as compiler
      compiler.Compiler().compile(passing_parameter, __file__ + '.tar.gz')
    

    【讨论】:

    • 很好的例子。谢谢
    猜你喜欢
    • 2021-12-01
    • 2021-04-25
    • 1970-01-01
    • 2017-01-09
    • 2020-03-11
    • 1970-01-01
    • 2020-06-12
    • 2020-10-26
    • 2018-07-08
    相关资源
    最近更新 更多