【问题标题】:Kryo encoder v.s. RowEncoder in Spark DatasetKryo 编码器与Spark 数据集中的行编码器
【发布时间】:2021-01-19 01:29:00
【问题描述】:

以下示例的目的是了解 Spark Dataset 中两种编码器的区别。

我可以这样做:

val df = Seq((1, "a"), (2, "d")).toDF("id", "value")

import org.apache.spark.sql.{Encoder, Encoders, Row}
import org.apache.spark.sql.catalyst.encoders.RowEncoder
import org.apache.spark.sql.types._

val myStructType = StructType(Seq(StructField("id", IntegerType), StructField("value", StringType)))
implicit val myRowEncoder = RowEncoder(myStructType)

val ds = df.map{case row => row}
ds.show

//+---+-----+
//| id|value|
//+---+-----+
//|  1|    a|
//|  2|    d|
//+---+-----+

我也可以这样做:

val df = Seq((1, "a"), (2, "d")).toDF("id", "value")

import org.apache.spark.sql.{Encoder, Encoders, Row}
import org.apache.spark.sql.catalyst.encoders.RowEncoder
import org.apache.spark.sql.types._

implicit val myKryoEncoder: Encoder[Row] = Encoders.kryo[Row] 

val ds = df.map{case row => row}
ds.show

//+--------------------+
//|               value|
//+--------------------+
//|[01 00 6F 72 67 2...|
//|[01 00 6F 72 67 2...|
//+--------------------+

代码的唯一区别是:一个使用Kryo编码器,另一个使用RowEncoder

问题:

  • 这两者有什么区别?
  • 为什么一个显示编码值,另一个显示人类可读值?
  • 我们什么时候应该使用哪个?

【问题讨论】:

  • 嗨@thebluephantom,感谢您的回答,但我也想知道前两个问题,甚至是第三个问题,Spark SQL 标准没有使用 Kryo,但是它仍然用于自定义类型,这就是我关心的问题
  • 我将在明天更新答案并运行一些示例。我注意到你的第二个例子不正确。

标签: apache-spark serialization apache-spark-sql apache-spark-dataset kryo


【解决方案1】:

Encoders.kryo 只是创建了一个编码器,它使用 Kryo序列化 T 类型的对象

RowEncoder 是 Scala 中的一个对象,具有 apply 和其他工厂方法。 RowEncoder 可以从模式创建 ExpressionEncoder[Row]。 在内部,apply 为 Row 类型创建一个 BoundReference,并为输入架构返回一个 ExpressionEncoder[Row]、一个 CreateNamedStruct 序列化器(使用 serializerFor 内部方法)、一个用于架构的反序列化器和 Row 类型

RowEncoder 了解架构并将其用于序列化和反序列化。

Kryo 比 Java 序列化(通常高达 10 倍)要快得多且更紧凑,但不支持所有 Serializable 类型,并且需要您提前注册将在程序中使用的类以获得最佳性能。

Kryo 适合高效存储大型数据集和网络密集型应用程序。

有关更多信息,您可以参考以下链接:

https://jaceklaskowski.gitbooks.io/mastering-spark-sql/content/spark-sql-RowEncoder.html https://jaceklaskowski.gitbooks.io/mastering-spark-sql/content/spark-sql-Encoders.html https://medium.com/@knoldus/kryo-serialization-in-spark-55b53667e7ab https://stackoverflow.com/questions/58946987/what-are-the-pros-and-cons-of-java-serialization-vs-kryo-serialization#:~:text=Kryo%20is%20significantly%20faster%20and,in%20advance%20for%20best%20performance.

【讨论】:

    【解决方案2】:

    根据 Spark 的文档,SparkSQL 不使用 Kryo 或 Java 序列化(标准)。

    Kryo 用于 RDD,而不是 Dataframe 或 DataSet。因此问题是 有点偏远了。

    Does Kryo help in SparkSQL? 这详细说明了自定义对象,但是...

    空闲时间后更新答案

    您的示例并不是我所说的自定义类型。他们是 只是带有原语的结构。没问题。

    Kryo 是一个序列化器,DS、DF 使用编码器来获得列优势。 Spark 在内部使用 Kryo 进行改组。

    这个用户定义的示例case class Foo(name: String, position: Point) 是我们可以使用 DS 或 DF 或通过 kryo 完成的示例。但是什么 Tungsten 和 Catalyst 合作的重点是“了解 数据结构”?从而能够优化。您还可以获得一个 使用 kryo 的单个二进制值,我发现了一些如何 成功地使用它,例如加入。

    KRYO 示例

    import org.apache.spark.sql.{Encoder, Encoders, SQLContext}
    import org.apache.spark.{SparkConf, SparkContext}
    import spark.implicits._
    
    case class Point(a: Int, b: Int)
    case class Foo(name: String, position: Point)
    
    implicit val PointEncoder: Encoder[Point] = Encoders.kryo[Point]
    implicit val FooEncoder: Encoder[Foo] = Encoders.kryo[Foo]
     
    val ds = Seq(new Foo("bar", new Point(0, 0))).toDS
    ds.show()
    

    返回:

    +--------------------+
    |               value|
    +--------------------+
    |[01 00 D2 02 6C 6...|
    +--------------------+
    

    使用案例类示例的 DS 编码器

    import org.apache.spark.sql.{Encoder, Encoders, SQLContext}
    import org.apache.spark.{SparkConf, SparkContext}
    import spark.implicits._
    
    case class Point(a: Int, b: Int)
    case class Foo(name: String, position: Point)
    
    val ds = Seq(new Foo("bar", new Point(0, 0))).toDS
    ds.show()
    

    返回:

    +----+--------+
    |name|position|
    +----+--------+
    | bar|  [0, 0]|
    +----+--------+
    

    这让我觉得 Spark、Tungsten、Catalyst 的最佳选择。

    现在,当涉及到 Any 时,会出现更复杂的情况,但 Any 不是一件好事:

    val data = Seq(
        ("sublime", Map(
          "good_song" -> "santeria",
          "bad_song" -> "doesn't exist")
        ),
        ("prince_royce", Map(
          "good_song" -> 4,
          "bad_song" -> "back it up")
        )
      )
    
    val schema = List(
        ("name", StringType, true),
        ("songs", MapType(StringType, StringType, true), true)
      )
    
    val rdd= spark.sparkContext.parallelize(data) 
    
    rdd.collect
    
    val df = spark.createDataFrame(rdd)
    df.show()
    df.printSchema()
    

    返回:

    Java.lang.UnsupportedOperationException: No Encoder found for Any.
    

    那么这个例子很有趣,它是一个有效的自定义对象用例 Spark No Encoder found for java.io.Serializable in Map[String, java.io.Serializable]。但我会远离这种情况。

    结论

    希望这会有所帮助。

    【讨论】:

      猜你喜欢
      • 2018-10-25
      • 2016-11-25
      • 2017-08-31
      • 2019-05-30
      • 2018-04-18
      • 1970-01-01
      • 2017-02-22
      • 2023-03-25
      • 2020-12-07
      相关资源
      最近更新 更多