在Spark中加入倾斜的数据集?

Raj*_*mar 23 join apache-spark

我正在使用Spark RDD加入两个大数据集.一个数据集非常偏斜,因此很少有执行程序任务需要很长时间才能完成工作.我该如何解决这个问题呢?

LiM*_*Bei 21

关于如何做的非常好的文章:https://datarus.wordpress.com/2015/05/04/fighting-the-skew-in-spark/

精简版:

  • 将随机元素添加到大型RDD并使用它创建新的连接键
  • 使用explode/flatMap将随机元素添加到小RDD以增加条目数并创建新的连接键
  • 在新的连接键上加入RDD,由于随机播种,现在将更好地分发

  • @kalpesh,还有一个参数需要考虑,Spark 中 shuffle 块的大小应该小于 2GB [SPARK-6235](https://issues.apache.org/jira/browse/SPARK-6235 )。我建议关注分区的大小,通常应该是 ~128MB,而不是分区数。我已经看到应用程序可以在许多分区(> 32k)下正常工作。 (2认同)

mor*_*007 15

假设您必须在A.id = B.id上加入两个表A和B. 让我们假设表A在id = 1时有偏差.

从A.id = B.id上的A连接B中选择A.id

解决倾斜连接问题有两种基本方法:

方法1:

将您的查询/数据集分成两部分 - 一部分仅包含倾斜,另一部分包含非倾斜数据.在上面的例子中.查询将成为 -

 1. select A.id from A join B on A.id = B.id where A.id <> 1;
 2. select A.id from A join B on A.id = B.id where A.id = 1 and B.id = 1;
Run Code Online (Sandbox Code Playgroud)

第一个查询不会有任何偏差,因此ResultStage的所有任务将在大致相同的时间完成.

如果我们假设B只有几行且B.id = 1,那么它将适合内存.因此,第二个查询将转换为广播连接.这在Hive中也称为Map-side join.

参考:https://cwiki.apache.org/confluence/display/Hive/Skewed+Join+Optimization

然后可以合并两个查询的部分结果以获得最终结果.

方法2:

上面的LeMuBei也提到过,第二种方法试图通过追加额外的列来随机化连接键.脚步:

  1. 在较大的表(A)中添加一列,例如skewLeft,并使用0到N-1之间的随机数填充所有行.

  2. 在较小的表(B)中添加一列,比如skewRight.复制较小的表N次.因此,对于每个原始数据副本,新skewRight列中的值将从0到N-1变化.为此,您可以使用explode sql/dataset运算符.

在1和2之后,加入2个数据集/表,并将连接条件更新为 -

                *A.id = B.id && A.skewLeft = B.skewRight*
Run Code Online (Sandbox Code Playgroud)

参考:https://datarus.wordpress.com/2015/05/04/fighting-the-skew-in-spark/


Jas*_*ans 12

根据您遇到的特定偏差类型,可能有不同的方法来解决它.基本思路是:

  • 修改您的连接列,或创建一个新的连接列,该列不会偏斜但仍保留足够的信息来进行连接
  • 在非倾斜列上进行连接 - 生成的分区不会偏斜
  • 在连接之后,您可以将连接列更新回您的首选格式,或者在创建新列时将其删除

如果倾斜的数据参与连接,那么LiMuBei的答案中引用的"反对火花中的歪斜"一文是一种很好的技巧.在我的例子中,偏差是由连接列中的大量空值引起的.空值未参与连接,但由于连接列上的Spark分区,后连接分区非常偏斜,因为有一个包含所有空值的巨大分区.

我通过添加一个新列来解决它,该列将所有空值更改为分布良好的临时值,例如"NULL_VALUE_X",其中X被替换为介于1和10,000之间的随机数,例如(在Java中):

// Before the join, create a join column with well-distributed temporary values for null swids.  This column
// will be dropped after the join.  We need to do this so the post-join partitions will be well-distributed,
// and not have a giant partition with all null swids.
String swidWithDistributedNulls = "swid_with_distributed_nulls";
int numNullValues = 10000; // Just use a number that will always be bigger than number of partitions
Column swidWithDistributedNullsCol =
    when(csDataset.col(CS_COL_SWID).isNull(), functions.concat(
        functions.lit("NULL_SWID_"),
        functions.round(functions.rand().multiply(numNullValues)))
    )
    .otherwise(csDataset.col(CS_COL_SWID));
csDataset = csDataset.withColumn(swidWithDistributedNulls, swidWithDistributedNullsCol);
Run Code Online (Sandbox Code Playgroud)

然后加入这个新列,然后加入:

outputDataset.drop(swidWithDistributedNullsCol);
Run Code Online (Sandbox Code Playgroud)


Rap*_*oth 0

您可以尝试将“倾斜”的 RDD 重新分区到更多分区,或者尝试增加分区spark.sql.shuffle.partitions(默认为 200)。

对于您的情况,我会尝试将分区数量设置为远高于执行程序的数量。

  • Spark.sql.shuffle.partitions 无助于数据倾斜。将有 200 个分区,但其中只有少数有数据。 (4认同)
  • @prakharjain 这并不完全正确。增加分区数量将减少具有许多记录的两个键被放置在同一分区中的机会。 (2认同)