【发布时间】:2017-05-24 00:26:56
【问题描述】:
我为我的 Spark 作业启用了 Kryo 序列化,启用了需要注册的设置,并确保我的所有类型都已注册。
val conf = new SparkConf()
conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
conf.set("spark.kryo.registrationRequired", "true")
conf.registerKryoClasses(classes)
conf.registerAvroSchemas(avroSchemas: _*)
该作业的挂钟时间性能下降了约 20%,洗牌的字节数增加了近 400%。
鉴于Spark documentation 建议 Kryo 应该更好,这让我感到非常惊讶。
Kryo 比 Java 序列化要快得多且更紧凑(通常高达 10 倍)
我在 Spark 的 org.apache.spark.serializer.KryoSerializer 和 org.apache.spark.serializer.JavaSerializer 实例上手动调用了 serialize 方法,并提供了我的数据示例。结果与 Spark 文档中的建议一致:Kryo 产生了 98 个字节; Java 产生了 993 个字节。这确实是 10 倍的改进。
一个可能令人困惑的因素是被序列化和洗牌的对象实现了 Avro GenericRecord 接口。我尝试在SparkConf 中注册 Avro 模式,但没有任何改善。
我尝试创建新类来对简单的 Scala case classes 数据进行洗牌,不包括任何 Avro 机器。它并没有提高 shuffle 性能或交换的字节数。
Spark 代码最终归结为以下内容:
case class A(
f1: Long,
f2: Option[Long],
f3: Int,
f4: Int,
f5: Option[String],
f6: Option[Int],
f7: Option[String],
f8: Option[Int],
f9: Option[Int],
f10: Option[Int],
f11: Option[Int],
f12: String,
f13: Option[Double],
f14: Option[Int],
f15: Option[Double],
f16: Option[Double],
f17: List[String],
f18: String) extends org.apache.avro.specific.SpecificRecordBase {
def get(f: Int) : AnyRef = ???
def put(f: Int, value: Any) : Unit = ???
def getSchema(): org.apache.avro.Schema = A.SCHEMA$
}
object A extends AnyRef with Serializable {
val SCHEMA$: org.apache.avro.Schema = ???
}
case class B(
f1: Long
f2: Long
f3: String
f4: String) extends org.apache.avro.specific.SpecificRecordBase {
def get(field$ : Int) : AnyRef = ???
def getSchema() : org.apache.avro.Schema = B.SCHEMA$
def put(field$ : Int, value : Any) : Unit = ???
}
object B extends AnyRef with Serializable {
val SCHEMA$ : org.apache.avro.Schema = ???
}
def join(as: RDD[A], bs: RDD[B]): (Iterable[A], Iterable[B]) = {
val joined = as.map(a => a.f1 -> a) cogroup bs.map(b => b.f1 -> b)
joined.map { case (_, asAndBs) => asAndBs }
}
您是否知道可能发生了什么,或者我如何才能获得 Kryo 应该提供的更好性能?
【问题讨论】:
-
你能发布示例案例类和工作吗?那么回答这个问题会容易得多
-
好点子,@T.Gawęd。用简化的代码更新。
-
你是如何衡量你的代码的?
-
@YuvalItzchakov 我根据每单位时间处理的记录数来衡量性能。我确保使用了相同数量的工人。我进行了很多试验。趋势很明显。我通过从 Spark UI 读取产生
cogroup输入的阶段的值来测量混洗字节。 -
你能确保你通过设置 sparkConf.set("spark.kryo.registrationRequired", "true") 注册了所有使用的东西吗?
标签: scala performance apache-spark avro kryo