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]

Reply via email to