【发布时间】:2020-03-14 04:30:53
【问题描述】:
我有一个数据框,它有两列 a 和 b,其中 b 列中的值是 a 列中值的子集。例如:
df
+---+---+
| a| b|
+---+---+
| 1| 2|
| 1| 3|
| 2| 1|
| 3| 2|
+---+---+
我想生成一个包含a 和anti_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