【问题标题】:How do I optimise non-equi-joins in Spark SQL? [duplicate]如何优化 Spark SQL 中的非等连接? [复制]
【发布时间】:2018-10-02 14:53:03
【问题描述】:

我有两个数据框需要使用具有两个连接谓词的非等连接(即不等连接)连接在一起。

一个数据框是一个直方图DataFrame[bin: bigint, lower_bound: double, upper_bound: double]
另一个数据框是观察集合DataFrame[id: bigint, observation: double]

我需要确定每个观察值属于直方图的哪个 bin,如下所示:

observations_df.join(histogram_df, 
    (
        (observations_df.observation >= histogram_df.lower_bound) &
        (observations_df.observation < histogram_df.upper_bound)
    )
   )

基本上它很慢,我正在寻找一些关于如何让它更快的建议。

下面是一些演示问题的示例代码。 observations_df 包含 100000 行,当 histogram_df 中的行数变得适当大(比如说number_of_bins = 500000)时,它变得非常非常慢,我确信这是因为我正在做一个非等值连接。如果您运行此代码,然后使用 number_of_rows 的值,从较低的值开始,然后增加,直到缓慢的性能很明显

from pyspark.sql.functions import lit, col, lead
from pyspark.sql.types import *
from pyspark.sql import SparkSession
from pyspark.sql.types import *
from pyspark.sql.functions import rand
from pyspark.sql import Window
spark = SparkSession \
    .builder \
    .getOrCreate()

number_of_bins = 500000

bin_width = 1.0 / number_of_bins
window = Window.orderBy('bin')
histogram_df = spark.range(0, number_of_bins)\
    .withColumnRenamed('id', 'bin')\
    .withColumn('lower_bound', 0 + lit(bin_width) * col('bin'))\
    .select('bin', 'lower_bound', lead('lower_bound', 1, 1.0).over(window).alias('upper_bound'))
observations_df = spark.range(0, 100000).withColumn('observation', rand())
observations_df.join(histogram_df, 
        (
            (observations_df.observation >= histogram_df.lower_bound) &
            (observations_df.observation < histogram_df.upper_bound)
        )
       ).groupBy('bin').count().head(15)

【问题讨论】:

标签: python apache-spark apache-spark-sql


【解决方案1】:

不建议将不等连接用于 spark join。通常,我会生成一个新列作为此类操作的连接键。 但是,对于您的情况,您不需要加入来确定每个观察值属于直方图的哪个 bin,因为可以预先计算每个 bin 的上限和下限,并且您可以使用观察值计算 bin。

您可以做的是编写一个 UDF,它会为您找到 bin 并将 bin 作为新列返回。 你可以参考pyspark: passing multiple dataframe fields to udf

【讨论】:

  • 感谢您的回复。在这个人为的演示中,是的,可以导出每个 bin 的上限和下限,但这只是一个演示。在我的现实世界中,情况并非如此。
  • 如果您坚持需要不等连接,那么您需要提供要连接的两个表的大小。调整将取决于尺寸。如果 histogram_df 很小(我假设),您可以缓存它并进行广播连接。如果 histogram_df 更小到可以放入字典(或其他数据结构)中,则可以广播字典并使用二进制搜索编写 UDF 来查找 bin。
  • 上面的demo代码特意代表了我的真实场景,实际上真实场景更糟糕,因为histogram_df中的number_of_bins实际上>1000000。因此,不幸的是广播连接不适合(我已经尝试过:))。 observations_df 在我的真实场景中有大约 100000 行,因此这就是我在上面的演示中填充的行数。
  • 我自己也在想这个问题,如果你做非等价连接,Catalyst 如何处理这个问题。
  • 但是如果你确实想要一个不相等的连接呢?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-10-30
  • 2012-06-28
相关资源
最近更新 更多