【问题标题】:Incorrect column names in spark datasetspark数据集中的列名不正确
【发布时间】:2016-09-19 12:05:54
【问题描述】:

我有下一个案例类:

case class Data[T](field1: String, field2: T)

我正在使用带有下一个隐式的 kryo 序列化程序:

implicit def single[A](implicit c: ClassTag[A]): Encoder[A] = Encoders.kryo[A](c)

implicit def tuple2[A1, A2](implicit e1: Encoder[A1], e2: Encoder[A2]): Encoder[(A1, A2)] =
        Encoders.tuple[A1, A2](e1, e2)

...

我尝试进行下一次加入:

val ds1 = someDataframe1.as[(String, T)].map(row => Data(row._1, row._2))
val ds1 = someDataframe2.as[(String, T)].map(row => Data(row._1, row._2))
ds1.joinWith(ds2, col("field1") === col("field1"), "left_outer")

之后我得到了下一个异常:

org.apache.spark.sql.AnalysisException: cannot resolve 'field1' given input columns: [value, value];

我的数据集中的列名发生了什么?

UPD: 当我打电话给ds1.schema 时,我得到了下一个输出:

StructField(name = value,dataType = BinaryType, nullable = true)

我认为我对 kryo 序列化有问题(没有架构元数据,案例类被序列化为没有名称的单个 blob 字段)。我还注意到,当T 是 kryo 已知类(Int,String)或案例类时,一切都很好。但是当T 是某个Java bean 时,我将我的数据数据集架构作为单个blob 未命名字段。

Spark 版本 1.6.1

【问题讨论】:

  • 我收到不同的错误消息(Spark 版本?)。无论如何,您是否尝试过col("_1.field1") === col("_2.field2) 作为加入条件?这种方式对我有用。
  • @Beryllium 请检查更新的问题

标签: scala apache-spark-sql kryo


【解决方案1】:

您刚刚创建了一个只有一列类型为 Data 的数据集,并且您无法访问该对象内的字段(fields1)。

dataframe: |value|
           |------
           |Data |
           |------
           |Data |
           |------

您可以尝试将您的数据框转换为数据集: val ds1 = someDataframe1.as[数据]

dataframe: |field1|field2|
           |------
           |String|  T   |
           |------
           |String|  T   |
           |------

或者如果你仍然想使用你的数据框,尝试改变搜索条件:

ds.joinWith(ds2, df.col("value").field1 === df2.col("value").field2)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-12-26
    • 2018-12-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-11-11
    相关资源
    最近更新 更多