PySpark Koalas下pandas三类默认分布式索引运行行为与预期不符咨询
PySpark Koalas三类默认索引测试疑问解答
核心理解偏差梳理
首先明确测试结果和官方说明不符的核心原因:
- 观测的样本仅来自
head()返回的第一个分区的前N行,无法代表全局索引的分布特征 - 官方说明里的分区、排序特征,是索引生成逻辑的特性,不是所有操作下都会立刻触发分区合并、数据重排
各索引类型问题对应解答
1. sequence类型索引分区数疑问
观测到分区数为8而非1,是因为官方说明中「数据会被收集到同一节点处理」的触发条件是:执行了需要按索引连续读取全量数据的操作(比如dfsp.to_pandas()、全局按索引切片)。仅生成sequence索引的阶段,只会做全局排序生成连续的索引值,不会修改原始读入的文件分区数,也不会主动把所有数据合并到单分区,避免不必要的性能损耗。
可以尝试执行dfsp.to_pandas()后再观测底层分区,会发现此时会触发全量数据拉取到单节点的逻辑。
2. distributed-sequence类型索引排序疑问
看到前5行索引和原始id顺序一致,是因为读入的CSV文件本身是按id从小到大顺序存储的,Spark读CSV时按文件偏移量拆分的8个分区,本身也是按原始数据顺序排列的。
distributed-sequence索引的生成逻辑是:先统计每个分区的行数,计算每个分区的索引起始值(上一个分区的最大索引+1),再在分区内生成连续的行号,拼出全局连续的索引。这种逻辑下如果原始分区是按顺序排列的,索引自然和原始数据顺序一致;如果原始数据是经过repartition打散的,索引就不会和原始id顺序对齐,符合官方「不保证结果自然排序」的说明。
可以用如下代码复现官方描述的行为:
ps.set_option('compute.default_index_type','distributed-sequence') df = ps.read_csv('df_s.csv') # 先把数据随机打散分区 df = df.spark.repartition(8) print(df.head())
此时输出的行id和索引值不再是连续对齐的。
3. distributed类型索引间隔疑问
看到前5行索引连续,同样是因为仅观测了第一个分区的内容。distributed索引的生成逻辑是:索引值 = 分区ID * 分区内最大行号阈值 + 分区内行号,这个阈值默认是非常大的数(通常远大于单个分区的行数),所以同一个分区内的索引是连续的,跨分区才会出现跳变。
可以用如下代码复现跳变特征:
ps.set_option('compute.default_index_type','distributed') dfsp = ps.read_csv('df_s.csv') # 取第二个分区的第一行数据 print(dfsp.filter(dfsp.id > 2500000).head(1))
此时输出的索引值会远大于2500000,出现明显的间隔。
三类默认索引的准确预期行为
| 索引类型 | 索引特征 | 性能开销 | 适用场景 |
|---|---|---|---|
| sequence | 全局连续、单调递增,和数据写入顺序完全一致 | 最高,需要全局排序shuffle | 需要索引和数据顺序严格对齐、数据量不大的场景 |
| distributed-sequence | 全局连续、单调递增,不保证和原始数据写入顺序对齐 | 中等,需要额外job统计每个分区的行数 | 需要索引全局连续、对顺序无要求的分布式场景 |
| distributed | 全局唯一、单调递增,跨分区存在固定间隔,不保证连续 | 最低,不需要shuffle,仅需要在分区内生成行号 | 仅需要索引作为唯一标识、不关注连续性和顺序的高性能场景 |
内容的提问来源于stack exchange,提问作者figs_and_nuts
相关产品推荐
相关产品推荐

