如何在Kubeflow组件间传递TensorFlow Dataset与张量?
解决Kubeflow组件间TensorFlow Dataset传递问题
问题根源
你遇到的AttributeError: 'Dataset' object has no attribute 'map'是因为:Kubeflow组件间传递的Dataset是Kubeflow的Artifact元数据对象,不是内存中的tf.data.Dataset实例。组件间传递的是数据集的存储路径引用,而非直接传递内存中的数据集对象,因此需要先将数据集持久化到存储,再在下游组件中加载。
解决方案步骤
1. 修改prepare_data组件:持久化TensorFlow Dataset到存储
在prepare_data函数内部,生成ratings、movies、train、test这几个tf.data.Dataset后,需要将它们保存到Kubeflow输出Dataset对应的存储路径(默认是GCS路径),使用TensorFlow的tf.data.experimental.save方法:
@component( packages_to_install = [ "pandas==1.3.4", "numpy==1.20.3", "unidecode", "nltk==3.6.5", "gcsfs==2023.1.0", "tensorflow==2.9.1" # 新增TF依赖,用于保存数据集 ], ) def prepare_data(dataset: str, data_artifact: Output[Dataset]) -> NamedTuple("Outputs", [("ratings", Dataset),("movies", Dataset),("train", Dataset),("test", Dataset)]): # 这里是你的数据处理逻辑,生成ratings_tfds、movies_tfds、train_tfds、test_tfds(tf.data.Dataset实例) # ... 你的数据处理代码 ... # 将每个tf Dataset保存到对应输出Artifact的路径 tf.data.experimental.save(ratings_tfds, ratings.path) tf.data.experimental.save(movies_tfds, movies.path) tf.data.experimental.save(train_tfds, train.path) tf.data.experimental.save(test_tfds, test.path) # 返回输出的NamedTuple实例 from collections import namedtuple Outputs = namedtuple("Outputs", ["ratings", "movies", "train", "test"]) return Outputs(ratings=ratings, movies=movies, train=train, test=test)
2. 修改train_model组件:从存储加载TensorFlow Dataset
在train_model函数中,通过Input[Dataset]的path属性获取数据集存储路径,再用tf.data.experimental.load加载为tf.data.Dataset实例,之后就可以正常调用map等方法:
@component( packages_to_install = [ "tensorflow-recommenders==0.7.0", "tensorflow==2.9.1", ], ) def train_model(epochs: int, ratings: Input[Dataset], movies: Input[Dataset], train: Input[Dataset], test: Input[Dataset], model_artifact: Output[Model]) -> NamedTuple("Outputs", [("model_artifact", Model)]): # 从存储路径加载tf.data.Dataset ratings_tfds = tf.data.experimental.load(ratings.path) movies_tfds = tf.data.experimental.load(movies.path) train_tfds = tf.data.experimental.load(train.path) test_tfds = tf.data.experimental.load(test.path) # 现在可以正常使用tf Dataset的方法 user_ids = ratings_tfds.map(lambda x: x["requisito"]) # ... 你的模型训练代码 ... # 保存模型到输出Artifact路径 tf.saved_model.save(model, model_artifact.path) # 返回输出 from collections import namedtuple Outputs = namedtuple("Outputs", ["model_artifact"]) return Outputs(model_artifact=model_artifact)
3. 管道代码无需修改
你的现有管道代码逻辑正确,组件间的输出引用传递没问题,只需确保上述组件内部的保存/加载逻辑正确即可。
关键注意事项
- 确保组件有访问GCS存储的权限:Vertex AI运行Kubeflow管道时,默认使用的服务账号需要有对应GCS桶的读写权限。
- TensorFlow版本兼容:
tf.data.experimental.save和load在TensorFlow 2.x版本中支持,你的2.9.1版本完全兼容。 - 数据集格式:保存和加载时,TensorFlow会自动处理数据集的结构,无需额外指定schema(如果数据集结构固定)。
内容的提问来源于stack exchange,提问作者razimbres
相关产品推荐
相关产品推荐

