【问题标题】:scala.MatchError during Spark 2.0.2 DataFrame unionSpark 2.0.2 DataFrame 联合期间的 scala.MatchError
【发布时间】:2017-06-03 10:52:28
【问题描述】:

我正在尝试使用 union 函数合并 2 个 DataFrame,一个包含旧数据,一个包含新数据。在我尝试向旧 DataFrame 动态添加新字段之前,这一直有效,因为我的架构正在发展。

这意味着我的旧数据将缺少一个字段,而新数据将拥有它。为了使联合工作,我正在使用下面的 EvolutionSchema 函数添​​加字段。

这导致了我粘贴在代码下方的输出/异常,包括我的调试打印。

列排序和使字段可以为空是试图通过使 DataFrame 尽可能相同来解决此问题,但它仍然存在。模式打印显示在这些操作之后它们看起来都是相同的。

任何有助于进一步调试的帮助将不胜感激。

import org.apache.spark.sql.functions.lit
import org.apache.spark.sql.types.{StructField, StructType}
import org.apache.spark.sql.{DataFrame, SQLContext}

object Merger {

  def apply(sqlContext: SQLContext, oldDataSet: Option[DataFrame], newEnrichments: Option[DataFrame]): Option[DataFrame] = {

    (oldDataSet, newEnrichments) match {
      case (None, None) => None
      case (None, _) => newEnrichments
      case (Some(existing), None) => Some(existing)
      case (Some(existing), Some(news)) => Some {

        val evolvedOldDataSet = evolveSchema(existing)

        println("EVOLVED OLD SCHEMA FIELD NAMES:" + evolvedOldDataSet.schema.fieldNames.mkString(","))
        println("NEW SCHEMA FIELD NAMES:" + news.schema.fieldNames.mkString(","))

        println("EVOLVED OLD SCHEMA FIELD TYPES:" + evolvedOldDataSet.schema.fields.map(_.dataType).mkString(","))
        println("NEW SCHEMA FIELD TYPES:" + news.schema.fields.map(_.dataType).mkString(","))

        println("OLD SCHEMA")
        existing.printSchema();
        println("PRINT EVOLVED OLD SCHEMA")
        evolvedOldDataSet.printSchema()
        println("PRINT NEW SCHEMA")
        news.printSchema()

        val nullableEvolvedOldDataSet = setNullableTrue(evolvedOldDataSet)
        val nullableNews = setNullableTrue(news)

        println("NULLABLE EVOLVED OLD")
        nullableEvolvedOldDataSet.printSchema()
        println("NULLABLE NEW")
        nullableNews.printSchema()

        val unionData =nullableEvolvedOldDataSet.union(nullableNews)

        val result = unionData.sort(
          unionData("timestamp").desc
        ).dropDuplicates(
          Seq("id")
        )
        result.cache()
      }
    }
  }

  def GENRE_FIELD : String = "station_genre"

  // Handle missing fields in old data
  def evolveSchema(oldDataSet: DataFrame): DataFrame = {
    if (!oldDataSet.schema.fieldNames.contains(GENRE_FIELD)) {

      val columnAdded = oldDataSet.withColumn(GENRE_FIELD, lit("N/A"))

      // Columns should be in the same order for union
      val columnNamesInOrder = Seq("id", "station_id", "station_name", "station_timezone", "station_genre", "publisher_id", "publisher_name", "group_id", "group_name", "timestamp")
      val reorderedColumns = columnAdded.select(columnNamesInOrder.head, columnNamesInOrder.tail: _*)

      reorderedColumns
    }
    else
      oldDataSet
  }

  def setNullableTrue(df: DataFrame) : DataFrame = {
    // get schema
    val schema = df.schema
    // create new schema with all fields nullable
    val newSchema = StructType(schema.map {
      case StructField(columnName, dataType, _, metaData) => StructField( columnName, dataType, nullable = true, metaData)
    })
    // apply new schema
    df.sqlContext.createDataFrame( df.rdd, newSchema )
  }

}

演变的旧模式字段名称: id,station_id,station_name,station_timezone,station_genre,publisher_id,publisher_name,group_id,group_name,timestamp

新的架构字段名称: id,station_id,station_name,station_timezone,station_genre,publisher_id,publisher_name,group_id,group_name,timestamp

进化的旧模式字段类型: StringType,LongType,StringType,StringType,StringType,LongType,StringType,LongType,StringType,LongType

新的模式字段类型: StringType,LongType,StringType,StringType,StringType,LongType,StringType,LongType,StringType,LongType

旧架构 根 |-- id: string (nullable = true) |-- station_id: long (nullable = true) |-- station_name: string (nullable = true) |-- station_timezone: string (nullable = true) |-- publisher_id: long (nullable = true) |-- publisher_name: string (nullable = true) |-- group_id: long (nullable = true) |-- group_name: string (nullable = true) |-- 时间戳:long (nullable = true)

打印 EVOLVED OLD SCHEMA 根 |-- id: string (nullable = true) |-- station_id: long (nullable = true) |-- station_name: string (nullable = true) |-- station_timezone: string (nullable = true) |-- station_genre: string (nullable = false) |-- publisher_id: long (nullable = true) |-- publisher_name: string (nullable = true) |-- group_id: long (nullable = true) |-- group_name: string (nullable = true) |-- 时间戳:long (nullable = true)

打印新的 SCHEMA 根 |-- id: string (nullable = true) |-- station_id: long (nullable = true) |-- station_name: string (nullable = true) |-- station_timezone: string (nullable = true) |-- station_genre: string (nullable = true) |-- publisher_id: long (nullable = true) |-- publisher_name: string (nullable = true) |-- group_id: long (nullable = true) |-- group_name: string (nullable = true) |-- 时间戳:long (nullable = true)

NULLABLE EVOLVED OLD root |-- id: string (nullable = true) |-- station_id: long (nullable = true) |-- station_name: string (nullable = true) |-- station_timezone: string (nullable = true) |-- station_genre: string (nullable = true) |-- publisher_id: long (nullable = true) |-- publisher_name: string (nullable = true) |-- group_id: long (nullable = true) |-- group_name: string (nullable = true) |-- 时间戳:long (nullable = true)

NULLABLE NEW root |-- id: string (nullable = true) |-- station_id: long (nullable = true) |-- station_name: string (nullable = true) |-- station_timezone: string (nullable = true) |-- station_genre: string (nullable = true) |-- publisher_id: long (nullable = true) |-- publisher_name: string (nullable = true) |-- group_id: long (nullable = true) |-- group_name: 字符串 (nullable = true) |-- 时间戳:long (nullable = true)

2017-01-18 15:59:32 错误 org.apache.spark.internal.Logging$class 执行者:91 - 阶段 2.0 (TID 4) 中的任务 1.0 中的异常 scala.MatchError: false (of class java.lang.Boolean) at org.apache.spark.sql.catalyst.CatalystTypeConverters$StringConverter$.toCatalystImpl(CatalystTypeConverters.scala:296) 在

...

com.companystuff.meta.uploader.Merger$.apply(Merger.scala:49)

...

原因:scala.MatchError: false (of class java.lang.Boolean) at org.apache.spark.sql.catalyst.CatalystTypeConverters$StringConverter$.toCatalystImpl(CatalystTypeConverters.scala:296) ...

【问题讨论】:

    标签: scala apache-spark spark-dataframe


    【解决方案1】:

    这是因为在实际数据中排序,即使其架构相同。 所以只需选择所有需要的列,然后进行联合查询。

    类似这样的:

    val columns:Seq[String]= ....
    val df = oldDf.select(columns:_*).union(newDf.select(columns:_*)
    

    希望对你有帮助

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2021-11-01
      • 2017-01-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多