Spark物理计划中ColumnarToRow算子输入批次数的确定逻辑问询
ColumnarToRow算子输入批次数的确定逻辑分析
我正在分析Spark查询生成的物理计划,该查询读取Parquet文件并执行聚合操作。物理计划里的ColumnarToRow算子带有「number of input batches」统计项,我想知道这个输入批次数是怎么确定的——它看起来和Parquet文件的行组数相关,但又不完全由行组数决定。
执行代码
df1 = spark.read.parquet('data/') .select('col1') .groupby('col1') .agg(f.count('col1').alias('ct')) .toPandas()
ColumnarToRow算子统计信息
ColumnarToRow number of output rows: 327,069 number of input batches: 80
输入批次数的确定逻辑
ColumnarToRow的输入批次数和Parquet行组有关,但还受以下核心因素影响:
- Parquet行组大小与Spark列批阈值:Spark通过
spark.sql.columnVector.batchSize配置列批的默认行数(默认10000行)。读取Parquet时,若单个行组的行数超过该阈值,会被拆分为多个列批;若行组行数小于等于阈值,则一个行组对应一个列批。 - 任务并行度与行组分配:Spark会根据Parquet文件的分区、大小等分配Task,每个Task处理若干行组。Task处理的行组数量(或拆分后的批数量)会直接计入ColumnarToRow的输入批次数。
- 数据过滤(若存在):如果查询包含谓词下推过滤,被过滤掉的行组不会生成列批;行组内数据过滤后剩余行数不足批大小的,也会形成一个独立的小批次。
回到你的场景:输入批次数为80,说明要么是处理了80个行数≤10000的Parquet行组,要么是部分行数超过10000的大行组被拆分后,总批次数达到80。你可以用parquet-tools meta data/命令查看Parquet文件的行组总数、每个行组的行数,就能直接验证上述逻辑。
内容的提问来源于stack exchange,提问作者Jason
相关产品推荐
相关产品推荐

