在由
year
和
month
用
spark.default.parallelism
设置为
4
,假设我要创建一个数据框架,由2017年的11~12个月和2018年的1~3个月两个来源组成
A
和
B
.
df = spark.read.parquet(
"A.parquet/_YEAR={2017}/_MONTH={11,12}",
"A.parquet/_YEAR={2018}/_MONTH={1,2,3}",
"B.parquet/_YEAR={2017}/_MONTH={11,12}",
"B.parquet/_YEAR={2018}/_MONTH={1,2,3}",
)
如果我得到了分区的数目,spark使用了
spark.default.parallelism并行性
默认情况下:
df.rdd.getNumPartitions()
Out[4]: 4
考虑到
df
我需要表演
join
和
groupBy
每个周期的操作,并且这些数据或多或少均匀地分布在每个周期上(每个周期大约1000万行):
问题
-
重新分区是否会提高后续操作的性能?
-
如果是,如果我有10个不同的时段(a和b都是每年5个),我是否应该按时段数重新分区并显式地引用要重新分区的列(
df.repartition(10,'_MONTH','_YEAR')
)?