根据指定黑名单标准的另一个DataFrame过滤Spark DataFrame

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)

  • 这就是为什么你不使用字符串,当你真的想要枚举时 - DataSet.join的scaladoc没有提到leftanti作为一个选项,所以不可能弄清楚,有什么选择 - 没有深入潜水.我开始觉得这需要一个Jira - 并且之前已经对API的选择感到不满. (7认同)
  • 值得庆幸的是,Jacek有完整的(或者我希望如此)希望加入的文档:https://jaceklaskowski.gitbooks.io/mastering-apache-spark/content/spark-sql-joins.html - 希望将此留在这里别人的生活更轻松. (6认同)
  • @RickMoritz此链接现在不可用. (3认同)
  • 说实话,当我写这个答案时,数据集是实验性的,我仍然不是一个大粉丝 (2认同)

Tag*_*gar 6

纯 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”的大小。