我想把每个分区都写成1GB的文件。因此,具有7gb数据的分区将作为7个文件写出,而具有0.3gb数据的分区将作为单个文件写出。
目前接受的答案在大多数情况下可能已经足够好了,但并不能满足将0.3gb分区写入单个文件的请求。相反,它会写出来
numPartitions
您要寻找的是一种根据数据分区的大小动态调整输出文件数量的方法。为此,我们将以10465355的使用方法为基础
rand()
repartition()
,并缩放
兰德()
基于该分区所需的文件数。
我将用Python提供一个演示,但是Scala中的方法基本相同。
from pyspark.sql import SparkSession
from pyspark.sql.functions import rand
spark = SparkSession.builder.getOrCreate()
skewed_data = (
spark.createDataFrame(
[(1,)] * 100 + [(2,)] * 10 + [(3,), (4,), (5,)],
schema=['id'],
)
)
partition_by_columns = ['id']
desired_rows_per_output_file = 10
partition_count = skewed_data.groupBy(partition_by_columns).count()
partition_balanced_data = (
skewed_data
.join(partition_count, on=partition_by_columns)
.withColumn(
'repartition_seed',
(
rand() * partition_count['count'] / desired_rows_per_output_file
).cast('int')
)
.repartition(*partition_by_columns, 'repartition_seed')
)
partition_count
. 如果您真的想动态地缩放每个分区的输出文件数,这是不可避免的。
from pyspark.sql.functions import spark_partition_id
(
skewed_data
.groupBy('id')
.count()
.orderBy('id')
.show()
)
(
partition_balanced_data
.select(
*partition_by_columns,
spark_partition_id().alias('partition_id'),
)
.groupBy(*partition_by_columns, 'partition_id')
.count()
.orderBy(*partition_by_columns, 'partition_id')
.show(30)
)
输出如下:
+---+-----+
| id|count|
+---+-----+
| 1| 100|
| 2| 10|
| 3| 1|
| 4| 1|
| 5| 1|
+---+-----+
+---+------------+-----+
| id|partition_id|count|
+---+------------+-----+
| 1| 7| 9|
| 1| 49| 6|
| 1| 53| 14|
| 1| 117| 12|
| 1| 126| 10|
| 1| 136| 11|
| 1| 147| 15|
| 1| 161| 7|
| 1| 177| 7|
| 1| 181| 9|
| 2| 85| 10|
| 3| 76| 1|
| 4| 197| 1|
| 5| 10| 1|
+---+------------+-----+
根据需要,每个输出文件大约有10行。
id=1
id=2
获取1个分区,然后
id={3,4,5}
这个解决方案平衡了输出文件的大小,而不考虑数据倾斜,并且不限制并行性
relying on
maxRecordsPerFile