【问题标题】:Does Kryo help in SparkSQL?Cryo 对 Spark SQL 有帮助吗?
【发布时间】:2018-03-14 06:15:58
【问题描述】:

Kryo 通过高效的序列化方法帮助提高 Spark 应用程序的性能。
我想知道,如果 Kryo 对 SparkSQL 有帮助,我应该如何使用它。
在 SparkSQL 应用程序中,我们会做很多基于列的操作,例如 df.select($"c1", $"c2"),而 DataFrame Row 的架构并不是完全静态的。
不确定如何为用例注册一个或多个序列化程序类。

例如:

case class Info(name: String, address: String)
...
val df = spark.sparkContext.textFile(args(0))
         .map(_.split(','))
         .filter(_.length >= 2)
         .map {e => Info(e(0), e(1))}
         .toDF
df.select($"name") ... // followed by subsequent analysis
df.select($"address") ... // followed by subsequent analysis

我认为为每个 select 定义案例类不是一个好主意。
或者如果我注册Info 像registerKryoClasses(Array(classOf[Info])) 一样有帮助

【问题讨论】:

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


【解决方案1】:

根据Spark's documentation,SparkSQL 不使用 Kryo 或 Java 序列化。

数据集与 RDD 类似,但是,它们不使用 Java 序列化或 Kryo,而是使用专门的编码器来序列化对象以通过网络进行处理或传输。虽然编码器和标准序列化都负责将对象转换为字节,但编码器是动态生成的代码,并且使用的格式允许 Spark 执行许多操作,例如过滤、排序和散列,而无需将字节反序列化回对象。

它们比 Java 或 Kryo 更轻量级,这是意料之中的(序列化是一项更可优化的工作,比如 3 个 long 和两个 int 的 Row),而不是类、它的版本描述、它的内部变量...)并且必须实例化它。

话虽如此,有一种方法可以使用 Kryo 作为编码器实现,请参见此处的示例:How to store custom objects in Dataset?。但这意味着将自定义对象(例如非产品类)存储在数据集中的解决方案,而不是专门针对标准数据帧。

如果没有 Java 序列化器的 Kryo,为自定义的非产品类创建编码器会受到一定的限制(请参阅关于用户定义类型的讨论),例如,从这里开始:Does Apache spark 2.2 supports user-defined type (UDT)?

【讨论】:

  • 这有改变吗?我可以使用隐式 val FooEncoder: Encoder[Foo] = Encoders.kryo[Foo] 创建 DS 或 DF 并获得二进制值。这是否意味着事情发生了变化?嗯,是的,因为上面的链接。
  • 我认为答案是正确的(目前)。数据集依赖于编码器,默认的(字符串、数字、案例和产品类...),不是“基于对象序列化的”,不使用 Kryo,并且速度快得多。如果您需要,Kryo 可以作为一个选项使用 - 通常会带来性能损失(在 CPU、堆内存和动态执行优化上更重)。人们经常提到 Kryo,因为早在 1.x 和 2.x 早期,RDD API 使用 java 序列化,Kryo 是一个推动力。然而,当前的 Dataset(/frame) 世界与那个时代大不相同。
  • 这也是我的结论。如何读取值等钨。谢谢
【解决方案2】:

您可以通过在您的 SparkConf 或通过 spark-submit 命令传递给 spark-submit 命令的自定义属性文件中将 spark.serializer 属性设置为 org.apache.spark.serializer.KryoSerializer 来将序列化程序设置为 kryo >--properties-file 标志。

当您配置 Kryo 序列化程序时,Spark 将在节点之间传输数据时透明地使用 Kryo。因此,您的 Spark SQL 语句应该会自动继承性能优势。

【讨论】:

  • 问题是,当我将类注册为conf.registerKryoClasses(Array(classOf[MyClass1])) 时,我应该如何定义案例类MyClass 字段,因为我感兴趣的列从select 更改为@ 987654326@?
  • 您能否提供更多关于MyClass 正在做什么的详细信息?我以为你只是想执行一个 SparkSQL 语句。
  • 此评论在使用 RDD 时是正确的,但在使用不使用 Kryo(也不是 java 序列化)的 Dataframe 时则不然。
猜你喜欢
  • 2018-07-09
  • 1970-01-01
  • 2011-08-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-03-07
相关资源
最近更新 更多