【问题标题】:Left anti join in groups左反加入团体
【发布时间】:2020-03-14 04:30:53
【问题描述】:

我有一个数据框,它有两列 ab,其中 b 列中的值是 a 列中值的子集。例如:

df
+---+---+
|  a|  b|
+---+---+
|  1|  2|
|  1|  3|
|  2|  1|
|  3|  2|
+---+---+

我想生成一个包含aanti_b 列的数据框,其中anti_b 列中的值是a 列中的任何值,例如a!=anti_b(a,anti_b) 行未出现在原始数据框中。所以在上面的数据框中,结果应该是:

anti df
+---+------+
|  a|anti_b|
+---+------+
|  3|     1|
|  2|     3|
+---+------+

这可以通过crossJoin 和对array_contains 的调用来完成,但它非常缓慢且效率低下。 有谁知道更好的 spark 成语来完成此任务,例如 anti_join

这是使用小数据框的低效示例,因此您可以看到我在追求什么:

df = spark.createDataFrame(pandas.DataFrame(numpy.array(
[[1,2],[1,3],[2,1],[3,2]]),columns=['a','b']))

crossed_df = df.select('a').withColumnRenamed('a','_a').distinct().crossJoin(df.select('a').withColumnRenamed('a','anti_b').distinct()).where(pyspark.sql.functions.col('_a')!=pyspark.sql.functions.col('anti_b'))
anti_df = df.groupBy(
  'a'
).agg(
  pyspark.sql.functions.collect_list('b').alias('bs')
).join(
  crossed_df,
  on=((pyspark.sql.functions.col('a')==pyspark.sql.functions.col('_a'))&(~pyspark.sql.functions.expr('array_contains(bs,anti_b)'))),
  how='inner'
).select(
  'a','anti_b'
)

print('df')
df.show()

print('anti df')
anti_df.show()

编辑:这也有效,但速度并不快:

df = spark.createDataFrame(pandas.DataFrame(numpy.array(
[[1,2],[1,3],[2,1],[3,2]]),columns=['a','b']))

crossed_df = df.select('a').distinct().crossJoin(df.select('a').withColumnRenamed('a','b').distinct()).where(pyspark.sql.functions.col('a')!=pyspark.sql.functions.col('b'))
anti_df = crossed_df.join(
  df,
  on=['a','b'],
  how='left_anti'
)

【问题讨论】:

  • 什么版本的 Spark?
  • @pault 2.4.1,但如果需要我可以升级到 2.4.4。
  • df.select('a').distinct() 的大小是多少?
  • @jxc 目前10万左右。

标签: apache-spark join pyspark apache-spark-sql


【解决方案1】:

这应该比你拥有的更好:

from pyspark.sql.functions import collect_set, expr

anti_df = df.groupBy("a").agg(collect_set("b").alias("bs")).alias("l")\
    .join(df.alias("r"), on=expr("NOT array_contains(l.bs, r.b)"))\
    .where("l.a != r.b")\
    .selectExpr("l.a", "r.b AS anti_b")\

anti_df.show()
#+---+------+
#|  a|anti_b|
#+---+------+
#|  3|     1|
#|  2|     3|
#+---+------+

如果您将此执行计划与您的方法进行比较,您会发现它更好(因为您可以将distinct 替换为collect_set),但它仍然具有笛卡尔积。

anti_df.explain()
#== Physical Plan ==
#*(3) Project [a#0, b#294 AS anti_b#308]
#+- CartesianProduct (NOT (a#0 = b#294) && NOT array_contains(bs#288, b#294))
#   :- *(1) Filter isnotnull(a#0)
#   :  +- ObjectHashAggregate(keys=[a#0], functions=[collect_set(b#1, 0, 0)])
#   :     +- Exchange hashpartitioning(a#0, 200)
#   :        +- ObjectHashAggregate(keys=[a#0], functions=[partial_collect_set(b#1, 0, 0)])
#   :           +- Scan ExistingRDD[a#0,b#1]
#   +- *(2) Project [b#294]
#      +- *(2) Filter isnotnull(b#294)
#         +- Scan ExistingRDD[a#293,b#294]

但是,如果没有更多信息,我认为没有任何方法可以避免针对这个特定问题的笛卡尔积。

【讨论】:

    猜你喜欢
    • 2017-08-28
    • 1970-01-01
    • 2018-08-13
    • 2023-03-21
    • 2011-03-14
    • 2010-10-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多