【发布时间】: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