【问题标题】:Spark SQL - Encoders for Tuple Containing a List or Array as an ElementSpark SQL - 包含列表或数组作为元素的元组的编码器
【发布时间】:2018-05-17 00:40:40
【问题描述】:

使用 Spark 2.2 + Java 1.8

我有两种自定义数据类型“Foo”和“Bar”。每个都实现了可序列化。'Foo' 与 'Bar' 具有一对多的关系,因此它们的关系表示为一个元组:

Tuple2<Foo, List<Bar>>

通常,当我有 1:1 的关系时,我可以像这样编码到我的自定义类型:

Encoder<Tuple2<Foo,Bar>> fooBarEncoder = Encoders.tuple(Encoders.bean(Foo.class),Encoders.bean(Bar.class));

然后用于编码我的数据集

Dataset<Tuple2<Foo,Bar>> fooBarSet = getSomeData().as(fooBarEncoder);

但是当我有一个列表(或数组)作为 Tuple2 元素时,我很难找到一种方法来为场景编码。我想做的是为第二个元素提供一个编码器,如下所示:

Encoder<Tuple2<Foo,List<Bar>>> fooBarEncoder = Encoders.tuple(Encoders.bean(Foo.class), List<Bar>.class);

然后编码到我的数据集:

Dataset<Tuple2<Foo,List<Bar>>> fooBarSet = getSomeData().as(fooBarEncoder)

但显然我不能在像 List 这样的参数化类型上调用 .class

我知道对于字符串和原始类型,spark 隐式支持数组,例如:

sparkSession.implicits().newStringArrayEncoder()

但是如何为自定义类类型的列表或数组创建编码器?

【问题讨论】:

    标签: java apache-spark apache-spark-sql spark-dataframe


    【解决方案1】:

    我不确定这种方法在您的设置中实现的效果如何,但这里可以。为您的列表创建一个包装类并尝试一下。

    public class BarList implements Serializable {
        List<Bar> list;
    
        public List<Bar> getList() {
            return list;
        }
        public void setList(List<Bar> l) {
            list = l;
        }
    }
    

    【讨论】:

      【解决方案2】:

      我不知道这是否可能。我尝试了以下 Scala,试图提供帮助,我想我可以通过首先教 spark 如何编码 X,然后是 List[X],最后是一个包含 List[X] 的元组(未在下面显示)来构建编码器:

      import org.apache.spark.sql.Encoders
      import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder
      import scala.beans.BeanProperty
      
      class X(@BeanProperty var field: String) extends Serializable
      case class Z(field: String)
      
      implicit val XEncoder1 = Encoders.bean(classOf[X])
      
      implicit val ZEncoder = Encoders.product[Z]
      
      val listXEncoder = ExpressionEncoder[List[X]] // doesn't work
      val listZEncoder = ExpressionEncoder[List[Z]]
      

      listZEncoder 工作正常

      切换使用

      implicit val XEncoder2 = org.apache.spark.sql.Encoders.kryo[X]
      

      仍然不适用于 listXEncoder

      错误最终出现在催化剂 ScalaReflection 中的某个位置,这超出了我的范围。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-03-26
        • 2020-10-20
        相关资源
        最近更新 更多