【发布时间】:2017-01-06 07:08:35
【问题描述】:
Dataframe A(百万条记录)的一列是create_date,modified_date
Dataframe B 500 记录有 start_date 和 end_date
目前的做法:
Select a.*,b.* from a join b on a.create_date between start_date and end_date
上面的查询在 sparkSQL 中执行笛卡尔积连接,它需要很长时间才能完成。 我可以通过其他方式实现相同的功能吗? 我尝试广播较小的 RDD
编辑:
spark version 1.4.1
No. of executors 2
Memmory/executor 5g
No. of cores 5
【问题讨论】:
-
"我尝试广播较小的 RDD" 发生了什么?
-
还是一样。 DAG 显示 shuffledHashJoin 和 cartesianProduct
-
内存表 A 的大小为 140MB,B 为 22kb,用于测试目的。 sparkUI 中此特定查询的输入显示为 2.2 GB
-
我怀疑 spark 中子句之间的 join 支持,因为 hive 只支持 join 上的相等性。
标签: scala apache-spark apache-spark-sql datastax-enterprise