Spark2.2使用cogroup关联KeyValueGroupedDataset的性能问题及替代方案
RDD fullOuterJoin的性能影响
RDD方案确实存在可观测的性能差异,核心原因如下:
- Dataset转RDD后会丢失Spark SQL层的Catalyst优化、Tungsten内存管理与代码生成能力,同时RDD采用的通用Java/Scala序列化开销远高于Dataset的专用序列化机制
- 数据规模在GB级及以上时,性能差距会比较明显,通常比同逻辑的Dataset实现慢2~3倍,小数据量下差异感知不明显
推荐的Dataset层面替代方案(兼容Spark 2.2版本)
方案1:保留cogroup写法,利用ds1 key唯一特性优化
现有cogroup写法完全可用,虽然对ds1执行groupByKey几乎无额外开销:因为ds1每个key仅对应1条数据,Spark内部对单条数据的分组会做特殊优化,不会产生多余的shuffle排序开销,仅需在处理逻辑中直接取ds1迭代器的首个元素即可:
case class Ds1(key: String, val1: String, val2: Int) case class Ds2(key: String, val1: String, val2: String) val ds1Grouped = df1.as[Ds1].groupByKey(_.key) val ds2Grouped = df2.as[Ds2].groupByKey(_.key) ds2Grouped.cogroup(ds1Grouped)( (k: String, ds2Iter: Iterator[Ds2], ds1Iter: Iterator[Ds1]) => { // ds1 key唯一,直接取头即可,无额外开销 val ds1Item = ds1Iter.toSeq.headOption // 你的自定义处理逻辑 } )
该方案shuffle次数与RDD方案一致,但全程享受Spark SQL层优化,性能远高于RDD实现。
方案2:先full outer join后分组处理
如果自定义逻辑可以基于同key的所有关联数据处理,可以先执行全外连接再按key分组,写法更直观:
val joinedDf = df2.join(df1, Seq("key"), "full_outer") joinedDf.groupByKey(row => row.getAs[String]("key")).mapGroups { (k, iter) => // 拆分ds1、ds2对应字段,执行自定义逻辑 }
方案3:ds1数据量较小时用广播变量优化
如果ds1数据量小于广播阈值(默认10M,可通过spark.sql.autoBroadcastJoinThreshold参数调整),可以直接将ds1广播为Map,完全避免shuffle,性能提升最明显:
// 广播ds1为<key, 数据>的Map val ds1Broadcast = spark.sparkContext.broadcast( df1.as[Ds1].map(item => (item.key, item)).collectAsMap() ) // 处理ds2中存在的key val ds2Result = df2.as[Ds2].groupByKey(_.key).mapGroups { (k, ds2Iter) => val ds1Item = ds1Broadcast.value.get(k) // 自定义逻辑 } // 若需要全外连接效果,补充处理仅ds1中存在的key val ds1OnlyKeys = df2.select("key").distinct().as[String] val ds1OnlyResult = df1.as[Ds1].filter(!ds1OnlyKeys.contains(_.key)).map { ds1Item => // 自定义处理仅ds1存在的key的逻辑 } // 合并结果 val finalResult = ds2Result.union(ds1OnlyResult)
内容的提问来源于stack exchange,提问作者gjin
相关产品推荐
相关产品推荐

