【问题标题】:UDF to extract only the file name from path in Spark SQLUDF 仅从 Spark SQL 中的路径中提取文件名
【发布时间】: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


【解决方案1】:

在 Scala 中:

#register udf
spark.udf
  .register("get_only_file_name", (fullPath: String) => fullPath.split("/").last)

#use the udf to get last token(filename) in full path
val initialDs = spark.read
  .option("dateFormat", conf.dateFormat)
  .schema(conf.schema)
  .csv(conf.path)
  .withColumn("input_file_name", get_only_file_name(input_file_name))

编辑:在 Java 中根据评论

#register udf
spark.udf()
  .register("get_only_file_name", (String fullPath) -> {
     int lastIndex = fullPath.lastIndexOf("/");
     return fullPath.substring(lastIndex, fullPath.length - 1);
    }, DataTypes.StringType);

import org.apache.spark.sql.functions.input_file_name    

#use the udf to get last token(filename) in full path
Dataset<Row> initialDs = spark.read()
  .option("dateFormat", conf.dateFormat)
  .schema(conf.schema)
  .csv(conf.path)
  .withColumn("input_file_name", get_only_file_name(input_file_name()));

【讨论】:

    【解决方案2】:

    借用一个相关问题here,下面的方法更便携,不需要自定义UDF。

    Spark SQL 代码片段: reverse(split(path, '/'))[0]

    Spark SQL 示例:

    WITH sample_data as (
    SELECT 'path/to/my/filename.txt' AS full_path
    )
    SELECT
          full_path
        , reverse(split(full_path, '/'))[0] as basename
    FROM sample_data
    

    解释: split() 函数将路径分成多个块,reverse() 将最后一项(文件名)放在数组前面,这样[0] 就可以只提取文件名。

    此处为完整代码示例:

      spark.sql(
        """
          |WITH sample_data as (
          |    SELECT 'path/to/my/filename.txt' AS full_path
          |  )
          |  SELECT
          |  full_path
          |  , reverse(split(full_path, '/'))[0] as basename
          |  FROM sample_data
          |""".stripMargin).show(false)
    
    

    结果:

    +-----------------------+------------+
    |full_path              |basename    |
    +-----------------------+------------+
    |path/to/my/filename.txt|filename.txt|
    +-----------------------+------------+
    

    【讨论】:

      【解决方案3】:

      commons io 是 spark 方式中自然/最简单的导入方式(无需添加额外的依赖...)

      import org.apache.commons.io.FilenameUtils
      
      getBaseName(String fileName)
      

      从完整文件名中获取基本名称,减去完整路径和扩展名。

      val baseNameOfFile = udf((longFilePath: String) => FilenameUtils.getBaseName(longFilePath))
      

      用法就像...

      yourdataframe.withColumn("shortpath" ,baseNameOfFile(yourdataframe("input_file_name")))
      .show(1000,false)
      

      【讨论】:

        猜你喜欢
        • 2010-10-01
        • 2016-06-19
        • 1970-01-01
        • 2020-11-20
        • 1970-01-01
        • 1970-01-01
        • 2018-04-21
        • 1970-01-01
        相关资源
        最近更新 更多