【发布时间】:2023-02-19 04:25:04
【问题描述】:
如何在 Pyspark 中使用加盐技术进行倾斜聚合。
假设我们有如下所示的倾斜数据,如何创建加盐列并将其用于聚合。
| city | state | count |
|---|---|---|
| Lachung | Sikkim | 3,000 |
| Rangpo | Sikkim | 50,000 |
| Gangtok | Sikkim | 3,00,000 |
| Bangalore | Karnataka | 2,50,00,000 |
| Mumbai | Maharashtra | 2,90,00,000 |
【问题讨论】:
如何在 Pyspark 中使用加盐技术进行倾斜聚合。
假设我们有如下所示的倾斜数据,如何创建加盐列并将其用于聚合。
| city | state | count |
|---|---|---|
| Lachung | Sikkim | 3,000 |
| Rangpo | Sikkim | 50,000 |
| Gangtok | Sikkim | 3,00,000 |
| Bangalore | Karnataka | 2,50,00,000 |
| Mumbai | Maharashtra | 2,90,00,000 |
【问题讨论】:
要在偏斜数据上使用加盐技术,我们需要创建一个名为“salt”的列。生成一个范围从 0 到 (spark.sql.shuffle.partitions - 1) 的随机编号。
表应如下所示,其中“salt”列的值介于 0 到 199 之间(因为在本例中分区大小为 200)。现在您可以对“城市”、“州”、“盐”使用 groupBy。
| city | state | salt |
|---|---|---|
| Lachung | Sikkim | 151 |
| Lachung | Sikkim | 102 |
| Lachung | Sikkim | 16 |
| Rangpo | Sikkim | 5 |
| Rangpo | Sikkim | 19 |
| Rangpo | Sikkim | 16 |
| Rangpo | Sikkim | 102 |
| Gangtok | Sikkim | 55 |
| Gangtok | Sikkim | 119 |
| Gangtok | Sikkim | 16 |
| Gangtok | Sikkim | 10 |
| Bangalore | Karnataka | 19 |
| Mumbai | Maharashtra | 0 |
| Bangalore | Karnataka | 199 |
| Mumbai | Maharashtra | 190 |
代码:
from pyspark.sql import SparkSession, functions as f
from pyspark.sql.types import (
StructType, StructField, IntegerType
)
salval = f.round(f.rand() * int(spark.conf.get("spark.sql.shuffle.partitions")) -1 )
record_df.withColumn("salt", f.lit(salval).cast(IntegerType()))
.groupBy("city", "state", "salt")
.agg(
f.count("city")
)
.drop("salt")
输出:
| city | state | count |
|---|---|---|
| Lachung | Sikkim | 3,000 |
| Rangpo | Sikkim | 50,000 |
| Gangtok | Sikkim | 3,00,000 |
| Bangalore | Karnataka | 2,50,00,000 |
| Mumbai | Maharashtra | 2,90,00,000 |
【讨论】: