【问题标题】:Dataframe transformations produce empty values数据框转换产生空值
【发布时间】:2020-10-22 07:41:52
【问题描述】:

我一直在尝试列出除元数据目录之外的其他目录中 Parquet 文件中的所有 Spark 数据帧。 目录结构如下:

dumped_data/
 - time=19424145
 - time=19424146
 - time=19424147
 - _spark_metadata

主要目标是避免从 _spark_metadata 目录读取数据。我创建了一个解决方案,但由于某种原因它不断返回空值。可能是什么原因?

解决办法如下:

 val dirNamesRegex: Regex = s"\\_spark\\_metadata*".r

def transformDf: Option[DataFrame] = {
 val filesDf = listPath(new Path(feedPath))(fsConfig)
      .map(_.getName)
      .filter(name => !dirNamesRegex.pattern.matcher(name).matches)
      .flatMap(path => sparkSession.parquet(Some(feedSchema))(path))

    if (!filesDf.isEmpty)
      Some(filesDf.reduce(_ union _))
    else None
 }

listPath - 在 hdfs 中列出数据文件的自定义方法。 feedSchema 是 StructType

没有 if on Some and None 我得到这个异常:

java.lang.UnsupportedOperationException: empty.reduceLeft
    at scala.collection.LinearSeqOptimized$class.reduceLeft(LinearSeqOptimized.scala:137)
    at scala.collection.immutable.List.reduceLeft(List.scala:84)
    at scala.collection.TraversableOnce$class.reduce(TraversableOnce.scala:208)
    at scala.collection.AbstractTraversable.reduce(Traversable.scala:104)

【问题讨论】:

    标签: regex scala apache-spark parquet


    【解决方案1】:

    在您的代码中,您有 3 个问题:

    1. 看来您可以使用== 运算符而不是正则表达式匹配。您知道要过滤的目录的具体名称,只需按名称过滤即可。
    2. 当我得到你的代码时,filesDf 类似于Traversable[DataFrame]。如果你想降低它的安全性,即使这个集合是空的,你也可以使用reduceLeftOption 而不是reduce
    3. 在您的transformDf 方法中,您尝试使用spark 过滤目录名称和读取数据,使用spark 进行调试也可能过于繁重。我建议您将您的逻辑分为两种不同的方法:首先 - 读取目录过滤 它们,其次 - 读取 数据和 union 它们成为一位将军DataFrame

    我提出这样的代码示例:

    不分逻辑的情况:

    def transformDf: Option[DataFrame] = {
      listPath(new Path(feedPath))(fsConfig)
        .map(_.getName)
        .filter(name => name == "_spark_metadata")
        .flatMap(path => sparkSession.parquet(Some(feedSchema))(path))
        .reduceLeftOption(_ union _)
    }
    

    逻辑分割案例:

    def getFilteredPaths: List[String] =
      listPath(new Path(feedPath))(fsConfig)
        .map(_.getName)
        .filter(name => name == "_spark_metadata")
    
    def transformDf: Option[DataFrame] = {
      getFilteredPaths
        .flatMap(path => sparkSession.parquet(Some(feedSchema))(path))
        .reduceLeftOption(_ union _)
    }
    

    在第二种方式中,您可以编写一些 轻量级 单元测试来调试您的路径提取,当您拥有正确的路径时,您可以轻松地从目录中读取数据并将它们合并。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-09-03
      • 2020-06-22
      • 1970-01-01
      • 2017-02-01
      • 1970-01-01
      • 1970-01-01
      • 2020-03-27
      • 1970-01-01
      相关资源
      最近更新 更多