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