【发布时间】:2017-04-12 10:24:41
【问题描述】:
Apache Spark 中有 input_file_name 函数,我使用它向 Dataset 添加具有当前正在处理的文件名称的新列。
问题是我想以某种方式自定义此函数以仅返回文件名,在 s3 上省略完整路径。
目前,我正在使用地图功能在第二步中替换路径:
val initialDs = spark.sqlContext.read
.option("dateFormat", conf.dateFormat)
.schema(conf.schema)
.csv(conf.path).withColumn("input_file_name", input_file_name)
...
...
def fromFile(fileName: String): String = {
val baseName: String = FilenameUtils.getBaseName(fileName)
val tmpFileName: String = baseName.substring(0, baseName.length - 8) //here is magic conversion ;)
this.valueOf(tmpFileName)
}
但我想使用类似的东西
val initialDs = spark.sqlContext.read
.option("dateFormat", conf.dateFormat)
.schema(conf.schema)
.csv(conf.path).withColumn("input_file_name", **customized_input_file_name_function**)
【问题讨论】:
-
.withColumn("input_file_name", get_only_file_name(input_file_name))。这里get_only_file_name是udf。
标签: java scala apache-spark apache-spark-sql spark-dataframe