This is an automated email from the ASF dual-hosted git repository.
zhouyuan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git
The following commit(s) were added to refs/heads/main by this push:
new e768fc81dd [VL][Iceberg] Propagate Parquet page row limit to Velox
(#12916)
e768fc81dd is described below
commit e768fc81dd3dd63ba1a99f5b88f6fe800ed1e8e5
Author: inf <[email protected]>
AuthorDate: Mon Sep 21 07:06:02 2026 +0000
[VL][Iceberg] Propagate Parquet page row limit to Velox (#12916)
---
.../execution/AbstractIcebergWriteExec.scala | 4 ++
.../execution/enhanced/VeloxIcebergSuite.scala | 83 ++++++++++++++++++++++
.../apache/gluten/execution/IcebergWriteExec.scala | 7 +-
3 files changed, 93 insertions(+), 1 deletion(-)
diff --git
a/backends-velox/src-iceberg/main/scala/org/apache/gluten/execution/AbstractIcebergWriteExec.scala
b/backends-velox/src-iceberg/main/scala/org/apache/gluten/execution/AbstractIcebergWriteExec.scala
index e9da0218d5..cfebdca24c 100644
---
a/backends-velox/src-iceberg/main/scala/org/apache/gluten/execution/AbstractIcebergWriteExec.scala
+++
b/backends-velox/src-iceberg/main/scala/org/apache/gluten/execution/AbstractIcebergWriteExec.scala
@@ -33,6 +33,9 @@ import scala.collection.JavaConverters._
abstract class AbstractIcebergWriteExec extends IcebergWriteExec {
+ private val parquetPageRowLimitSession =
+ "spark.gluten.sql.columnar.backend.velox.parquet_writer_page_row_limit"
+
// the writer factory works for both batch and streaming
private def createIcebergDataWriteFactory(schema: StructType):
IcebergDataWriteFactory = {
val writeSchema = IcebergWriteUtil.getWriteSchema(write)
@@ -50,6 +53,7 @@ abstract class AbstractIcebergWriteExec extends
IcebergWriteExec {
Seq(
PARQUET_PAGE_SIZE_BYTES.key -> getParquetPageSizeBytes,
COLUMNAR_PARQUET_WRITE_BLOCK_SIZE.key -> getParquetRowGroupSizeBytes,
+ parquetPageRowLimitSession -> getParquetPageRowLimit,
MAX_TARGET_FILE_SIZE_SESSION.key -> getTargetFileSizeBytes,
PARQUET_DICT_SIZE_BYTES.key -> getDictSizeBytes
).foreach {
diff --git
a/backends-velox/src-iceberg/test/scala/org/apache/gluten/execution/enhanced/VeloxIcebergSuite.scala
b/backends-velox/src-iceberg/test/scala/org/apache/gluten/execution/enhanced/VeloxIcebergSuite.scala
index 00c6e54b58..3ad8c7ffbc 100644
---
a/backends-velox/src-iceberg/test/scala/org/apache/gluten/execution/enhanced/VeloxIcebergSuite.scala
+++
b/backends-velox/src-iceberg/test/scala/org/apache/gluten/execution/enhanced/VeloxIcebergSuite.scala
@@ -656,6 +656,89 @@ class VeloxIcebergSuite extends IcebergSuite {
}
}
+ test("iceberg table page row limit") {
+ val table = "iceberg_page_row_limit"
+
+ def dataPageRowCounts(table: String, columnName: String): Seq[Int] = {
+ val conf = spark.sparkContext.hadoopConfiguration
+ val files = spark.sql(s"""
+ SELECT file_path
+ FROM default.$table.files
+ """).collect().map(_.getString(0)).toSeq
+
+ files.flatMap {
+ file =>
+ val inputFile = HadoopInputFile.fromPath(new Path(file), conf)
+ val reader = ParquetFileReader.open(inputFile,
ParquetReadOptions.builder().build())
+
+ try {
+ val column = reader
+ .getFooter
+ .getFileMetaData
+ .getSchema
+ .getColumns
+ .asScala
+ .find(_.getPath.toSeq == Seq(columnName))
+ .getOrElse(fail(s"Column $columnName was not found in Parquet
file $file"))
+
+ val rowCounts = scala.collection.mutable.ArrayBuffer.empty[Int]
+ var rowGroup = reader.readNextRowGroup()
+ while (rowGroup != null) {
+ val pageReader = rowGroup.getPageReader(column)
+ pageReader.readDictionaryPage()
+
+ var page = pageReader.readPage()
+ while (page != null) {
+ rowCounts += page.getValueCount
+ page = pageReader.readPage()
+ }
+
+ rowGroup = reader.readNextRowGroup()
+ }
+ rowCounts
+ } finally {
+ reader.close()
+ }
+ }
+ }
+
+ withSQLConf("spark.sql.shuffle.partitions" -> "1") {
+ withTable(table) {
+ spark.sql(s"""
+ CREATE TABLE $table (
+ value SMALLINT
+ ) USING iceberg
+ TBLPROPERTIES (
+ 'write.format.default' = 'parquet',
+ 'write.parquet.compression-codec' = 'uncompressed',
+ 'write.parquet.page-size-bytes' = '1MB',
+ 'write.parquet.page-row-limit' = '1000'
+ )
+ """)
+
+ val df = spark.sql(s"""
+ INSERT INTO $table
+ SELECT CAST(id AS SMALLINT)
+ FROM range(0, 5000, 1, 1)
+ """)
+
+ assert(
+ df.queryExecution.executedPlan
+ .asInstanceOf[CommandResultExec]
+ .commandPhysicalPlan
+ .isInstanceOf[VeloxIcebergAppendDataExec])
+
+ val pageRowCounts = dataPageRowCounts(table, "value")
+ assert(
+ pageRowCounts.size > 1,
+ s"Expected the Iceberg page-row limit to create multiple data pages:
$pageRowCounts")
+ assert(
+ pageRowCounts.sum == 5000,
+ s"Expected 5000 values across all data pages: $pageRowCounts")
+ }
+ }
+ }
+
test("iceberg table row group size and Gluten override") {
val table = "iceberg_row_group_size"
val overrideTable = "iceberg_row_group_size_override"
diff --git
a/gluten-iceberg/src/main/scala/org/apache/gluten/execution/IcebergWriteExec.scala
b/gluten-iceberg/src/main/scala/org/apache/gluten/execution/IcebergWriteExec.scala
index d97aac4bd4..99638c064c 100644
---
a/gluten-iceberg/src/main/scala/org/apache/gluten/execution/IcebergWriteExec.scala
+++
b/gluten-iceberg/src/main/scala/org/apache/gluten/execution/IcebergWriteExec.scala
@@ -19,7 +19,7 @@ package org.apache.gluten.execution
import org.apache.gluten.backendsapi.BackendsApiManager
import org.apache.iceberg.{FileFormat, PartitionField, PartitionSpec, Schema,
TableProperties}
-import org.apache.iceberg.TableProperties.{ORC_COMPRESSION,
ORC_COMPRESSION_DEFAULT, PARQUET_COMPRESSION, PARQUET_COMPRESSION_DEFAULT,
PARQUET_DICT_SIZE_BYTES, PARQUET_DICT_SIZE_BYTES_DEFAULT,
PARQUET_PAGE_SIZE_BYTES, PARQUET_PAGE_SIZE_BYTES_DEFAULT,
PARQUET_ROW_GROUP_SIZE_BYTES, PARQUET_ROW_GROUP_SIZE_BYTES_DEFAULT}
+import org.apache.iceberg.TableProperties.{ORC_COMPRESSION,
ORC_COMPRESSION_DEFAULT, PARQUET_COMPRESSION, PARQUET_COMPRESSION_DEFAULT,
PARQUET_DICT_SIZE_BYTES, PARQUET_DICT_SIZE_BYTES_DEFAULT,
PARQUET_PAGE_ROW_LIMIT, PARQUET_PAGE_ROW_LIMIT_DEFAULT,
PARQUET_PAGE_SIZE_BYTES, PARQUET_PAGE_SIZE_BYTES_DEFAULT,
PARQUET_ROW_GROUP_SIZE_BYTES, PARQUET_ROW_GROUP_SIZE_BYTES_DEFAULT}
import org.apache.iceberg.avro.AvroSchemaUtil
import org.apache.iceberg.spark.source.IcebergWriteUtil
import org.apache.iceberg.types.Type.TypeID
@@ -56,6 +56,11 @@ trait IcebergWriteExec extends ColumnarV2TableWriteExec {
normalizeCapacityString(PARQUET_PAGE_SIZE_BYTES_DEFAULT.toString))
}
+ protected def getParquetPageRowLimit: String = {
+ val tableProps = IcebergWriteUtil.getTable(write).properties()
+ tableProps.getOrDefault(PARQUET_PAGE_ROW_LIMIT,
PARQUET_PAGE_ROW_LIMIT_DEFAULT.toString)
+ }
+
protected def getTargetFileSizeBytes: String = {
IcebergWriteUtil.getWriteConf(write).targetDataFileSize().toString
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]