如何在代码仓库的transform转换内获取目标数据集名称
解决方案
在平台代码仓库编写Python转换任务时,@transform_df默认会将声明的Input直接加载为DataFrame传入计算函数,不会携带数据集名称、路径等元信息。运行时输入数据集会被挂载到随机临时目录,直接通过os.listdir()遍历文件系统无法拿到Input和实际数据集的映射关系,因此无法通过路径遍历的方式提取数据集名称。
具体实现
- 换用
@transform装饰器替代@transform_df,此时传入计算函数的是Input对象本身,既可以读取数据集内容,也可以访问其元信息。 - 通过Input对象的
path属性获取数据集在平台的完整逻辑路径,截取路径最后一段即为数据集名称,该值不受临时挂载路径影响,始终固定准确。 - 给每个输入数据集对应的DataFrame新增存储来源名称的列,合并后调整列顺序,将名称列放在第一列,最后写出结果即可。
代码示例
from transforms.api import transform, Input, Output import pyspark.sql.functions as F @transform( output=Output("/folder/folder1/datasets/mydataset"), input_a=Input("A"), input_b=Input("B"), ) def compute(input_a, input_b, output): # 提取每个输入的数据集名称 name_a = input_a.path.split("/")[-1] name_b = input_b.path.split("/")[-1] # 读取数据集并新增来源列 df_a = input_a.dataframe().withColumn("dataset_name", F.lit(name_a)) df_b = input_b.dataframe().withColumn("dataset_name", F.lit(name_b)) # 合并数据集,调整列顺序将名称列放到第一列 merged_df = df_a.unionByName(df_b) final_columns = ["dataset_name"] + [col for col in merged_df.columns if col != "dataset_name"] final_df = merged_df.select(*final_columns) # 写出结果 output.write_dataframe(final_df)
原代码存在两处可修正的问题:
- 装饰器仅声明了2个输入,计算函数却定义了
df3参数,运行时会触发参数不匹配报错。- 输出路径中的
mydatset存在拼写错误,建议核对路径后再提交任务。
内容的提问来源于stack exchange,提问作者slowly study
相关产品推荐
相关产品推荐

