【问题标题】:Why is Spark broadcast exchange data size bigger than raw size on join?为什么 Spark 广播交换数据大小大于连接时的原始大小?
【发布时间】:2020-01-21 20:08:13
【问题描述】:

我正在对两个表 A 和 B 进行广播连接。 B 是使用以下 Spark SQL 创建的缓存表:

create table B as select segment_ids_hash from  stb_ranker.c3po_segments
      where
        from_unixtime(unix_timestamp(string(dayid), 'yyyyMMdd')) >= CAST('2019-07-31 00:00:00.000000000' AS TIMESTAMP)
      and
        segmentid_check('(6|8|10|12|14|371|372|373|374|375|376|582|583|585|586|587|589|591|592|594|596|597|599|601|602|604|606|607|609|610|611|613|615|616)', seg_ids) = true
cache table B

“segment_ids_hash”列是整数类型,结果包含 3640 万条记录。 缓存表大小约为140 MB,如下图

然后我按如下方式进行连接:

select count(*) from A broadcast join B on A.segment_ids_hash = B.segment_ids_hash

这里广播交换数据大小约为 3.2 GB。

我的问题是为什么广播交换数据大小 (3.2GB) 比原始数据大小 (~140 MB) 大得多。什么是间接费用?有什么办法可以减少广播交换数据的大小?

谢谢

【问题讨论】:

  • 你的集群大小是多少,你使用的是什么序列化,请更新成问题,当 Brodcast 数据对 100 个 140 MB 变量的引用时,它应该是 140 GB
  • Spark 广播将数据收集到驱动程序然后分派给每个执行程序,显示的大小是我认为通过网络发送的总数
  • @sramalingam24 我用不同数量的执行器进行了测试,广播字节大小没有改变。
  • 你知道直播的是哪一部吗?您可以查看生成的查询计划,其他的分区数也会起作用

标签: apache-spark apache-spark-sql


【解决方案1】:

Tl;博士:我也在学习数据大小指标的来源。这可能只是操作的估计大小,它可能无法反映数据的实际大小。暂时不用太担心。

完整版:

更新:回来纠正一些错误。我看到之前的答案缺乏深度,所以我会尽量深入挖掘(我对回答问题还是比较陌生)。

更新 2:改写,删除了一些过头的笑话(sry)

好的,所以这件事可能很长,但我认为这个指标并不是数据的直接大小。

首先,我对此进行了测试运行,以使用 200 个执行器和 4 个核心重现结果:

这返回了以下结果:

现在我看到了一些有趣的东西,因为我测试的 dataSize 大约是 1.2GB 而不是 3.2GB,这导致我阅读了 Spark 的源代码。

我去github的时候看到BroadcastExchange里面的4个数字对应这个: 第一个链接:BroadcastHashJoinExec:https://github.com/apache/spark/blob/master/sql/core/src/main/scala/org/apache/spark/sql/execution/exchange/BroadcastExchangeExec.scala

这个对应的数据大小:

我发现这里的关系 val 似乎是一个 HashedRelationBroadcastMode。

转到 HashedRelation https://github.com/apache/spark/blob/master/sql/core/src/main/scala/org/apache/spark/sql/execution/joins/HashedRelation.scala:

因为我们有 Some(Numrows)(它是 DF 的行数)。匹配用例用例一(第 926:927 行)

回到 HashedRelation 构造函数-y 部分:

由于join是hashed int,类型不是Long => join使用UnsafeHashedRelation

到 UnsafeHashedRelation:

现在我们去 UnsafeHashedRelation 中确定估计大小的地方,我发现了这个:

关注估计的大小,我们的目标是binaryMap对象(后面代码assign map = binaryMap)

然后到这里:

binaryMap是一个BytestoBytesMap,这里对应https://github.com/apache/spark/blob/master/core/src/main/java/org/apache/spark/unsafe/map/BytesToBytesMap.java

跳转到getTotalMemoryConsumption方法(获取estimateSize的那个),我们得到:

这是我目前的死胡同。只是我的两分钱,我不认为这是一个错误,而只是连接的估计大小,并且由于这是估计的大小,我并不认为它必须非常准确(是的,但它很奇怪在这种情况下是诚实的,因为差异很大)。

如果您想继续使用这个数据大小。一种方法是通过修改其构造函数的输入来直接影响 binaryMap 对象。回头看看:

有两个变量可以配置,MEMORY_OFFHEAP_ENABLED和BUFFER_PAGE大小。也许您可以在 spark-submit 期间尝试使用这两种配置。这也是为什么即使您更改了执行器和核心的数量,BroadcastExec 的大小也不会改变的原因。

因此,总而言之,我认为数据大小是由某种有趣的机制生成的估计值(我也在等待更专业的人来解释这一点,因为我正在研究它),而不是直接的大小在第一张图片(140 MB)中提到过。因此,可能不值得花太多时间来减少此特定指标的开销。

一些奖金相关的东西:

https://jaceklaskowski.gitbooks.io/mastering-spark-sql/spark-sql-SparkPlan-BroadcastExchangeExec.html

https://jaceklaskowski.gitbooks.io/mastering-spark-sql/spark-sql-UnsafeRow.html

【讨论】:

  • 感谢您的努力,但无论是“别想太多” 还是您显示的公式都不能为估算提供令人满意的解释,抱歉。用较小的数据集(每个 2 个元素)尝试了您的示例代码,data size 仍然超过 1MB。对我来说,这应该被认为是一个错误或正在发生的其他事情。
猜你喜欢
  • 2019-05-12
  • 2012-01-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-03-21
  • 1970-01-01
  • 2020-11-13
相关资源
最近更新 更多