You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 05:40:26