Pow*_*ers 24 dataframe apache-spark apache-spark-sql
我有一个largeDataFrame(多列和数十亿行)和一个smallDataFrame(单列和10,000行).
我想所有的行从过滤largeDataFrame每当some_identifier列在largeDataFrame比赛中的行之一smallDataFrame.
这是一个例子:
largeDataFrame
some_idenfitier,first_name
111,bob
123,phil
222,mary
456,sue
Run Code Online (Sandbox Code Playgroud)
smallDataFrame
some_identifier
123
456
Run Code Online (Sandbox Code Playgroud)
desiredOutput
111,bob
222,mary
Run Code Online (Sandbox Code Playgroud)
这是我丑陋的解决方案.
val smallDataFrame2 = smallDataFrame.withColumn("is_bad", lit("bad_row"))
val desiredOutput = largeDataFrame.join(broadcast(smallDataFrame2), Seq("some_identifier"), "left").filter($"is_bad".isNull).drop("is_bad")
Run Code Online (Sandbox Code Playgroud)
有更清洁的解决方案吗?
eli*_*sah 59
left_anti在这种情况下,您需要使用连接.
在左边的抗加入是一个相反的左半加入.
它根据给定的密钥从左表中的右表中过滤掉数据:
largeDataFrame
.join(smallDataFrame, Seq("some_identifier"),"left_anti")
.show
// +---------------+----------+
// |some_identifier|first_name|
// +---------------+----------+
// | 222| mary|
// | 111| bob|
// +---------------+----------+
Run Code Online (Sandbox Code Playgroud)
纯 Spark SQL 中的一个版本(以 PySpark 为例,但有一些小的变化同样适用于 Scala API):
def string_to_dataframe (df_name, csv_string):
rdd = spark.sparkContext.parallelize(csv_string.split("\n"))
df = spark.read.option('header', 'true').option('inferSchema','true').csv(rdd)
df.registerTempTable(df_name)
string_to_dataframe("largeDataFrame", '''some_identifier,first_name
111,bob
123,phil
222,mary
456,sue''')
string_to_dataframe("smallDataFrame", '''some_identifier
123
456
''')
anti_join_df = spark.sql("""
select *
from largeDataFrame L
where NOT EXISTS (
select 1 from smallDataFrame S
WHERE L.some_identifier = S.some_identifier
)
""")
print(anti_join_df.take(10))
anti_join_df.explain()
Run Code Online (Sandbox Code Playgroud)
将按预期输出 mary 和 bob:
[行(some_identifier=222, first_name='mary'),
行(some_identifier=111, first_name='bob')]
并且物理执行计划将显示它正在使用
== Physical Plan ==
SortMergeJoin [some_identifier#252], [some_identifier#264], LeftAnti
:- *(1) Sort [some_identifier#252 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(some_identifier#252, 200)
: +- Scan ExistingRDD[some_identifier#252,first_name#253]
+- *(3) Sort [some_identifier#264 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(some_identifier#264, 200)
+- *(2) Project [some_identifier#264]
+- Scan ExistingRDD[some_identifier#264]
Run Code Online (Sandbox Code Playgroud)
通知Sort Merge Join对于加入/反加入大约相同大小的数据集更有效。由于您已经提到小数据帧较小,因此我们应该确保 Spark 优化器选择Broadcast Hash Join在这种情况下效率更高的方法:
为此,我将更NOT EXISTS改为NOT IN条款:
== Physical Plan ==
SortMergeJoin [some_identifier#252], [some_identifier#264], LeftAnti
:- *(1) Sort [some_identifier#252 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(some_identifier#252, 200)
: +- Scan ExistingRDD[some_identifier#252,first_name#253]
+- *(3) Sort [some_identifier#264 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(some_identifier#264, 200)
+- *(2) Project [some_identifier#264]
+- Scan ExistingRDD[some_identifier#264]
Run Code Online (Sandbox Code Playgroud)
让我们看看它给了我们什么:
== Physical Plan ==
BroadcastNestedLoopJoin BuildRight, LeftAnti, ((some_identifier#302 = some_identifier#314) || isnull((some_identifier#302 = some_identifier#314)))
:- Scan ExistingRDD[some_identifier#302,first_name#303]
+- BroadcastExchange IdentityBroadcastMode
+- Scan ExistingRDD[some_identifier#314]
Run Code Online (Sandbox Code Playgroud)
请注意,Spark Optimizer 实际上选择了Broadcast Nested Loop Join而不是Broadcast Hash Join。前者没问题,因为我们只有两条记录要从左侧排除。
还要注意,两个执行计划都有,LeftAnti所以它类似于@eliasah 答案,但使用纯 SQL 实现。此外,它还表明您可以更好地控制物理执行计划。
附注。还要记住,如果右侧的数据框比左侧的数据框小得多,但又大于几条记录,则您确实希望拥有Broadcast Hash Joinand not Broadcast Nested Loop Joinnor Sort Merge Join。如果这没有发生,您可能需要调整spark.sql.autoBroadcastJoinThreshold因为它默认为 10Mb,但它必须大于“smallDataFrame”的大小。
| 归档时间: |
|
| 查看次数: |
19463 次 |
| 最近记录: |