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 24340453ec [VL][Iceberg] Fix target file size bytes to session config 
(#12901)
24340453ec is described below

commit 24340453ecf18d592de2a54016d1acb850cdae45
Author: inf <[email protected]>
AuthorDate: Mon Sep 28 08:50:48 2026 +0000

    [VL][Iceberg] Fix target file size bytes to session config (#12901)
---
 .../execution/AbstractIcebergWriteExec.scala       |   2 +-
 .../spark/source/IcebergWriteExecCapacityTest.java |  94 ++++++++++++
 .../execution/enhanced/VeloxIcebergSuite.scala     | 160 +++++++++++----------
 cpp/velox/utils/ConfigExtractor.cc                 |   2 +-
 .../apache/gluten/execution/IcebergWriteExec.scala |  12 +-
 5 files changed, 189 insertions(+), 81 deletions(-)

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 cfebdca24c..6df3b105f5 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
@@ -61,7 +61,7 @@ abstract class AbstractIcebergWriteExec extends 
IcebergWriteExec {
         val overrideValue = SQLConf.get.getConfString(key, null)
         if (overrideValue == null) {
           icebergProperties.put(key, value)
-        } else if (key == COLUMNAR_PARQUET_WRITE_BLOCK_SIZE.key) {
+        } else if (key != parquetPageRowLimitSession) {
           icebergProperties.put(key, normalizeCapacityString(overrideValue))
         }
     }
diff --git 
a/backends-velox/src-iceberg/test/java/org/apache/iceberg/spark/source/IcebergWriteExecCapacityTest.java
 
b/backends-velox/src-iceberg/test/java/org/apache/iceberg/spark/source/IcebergWriteExecCapacityTest.java
new file mode 100644
index 0000000000..e4255b725e
--- /dev/null
+++ 
b/backends-velox/src-iceberg/test/java/org/apache/iceberg/spark/source/IcebergWriteExecCapacityTest.java
@@ -0,0 +1,94 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.iceberg.spark.source;
+
+import org.apache.gluten.execution.IcebergWriteExec;
+
+import org.apache.iceberg.Table;
+import org.apache.iceberg.TableProperties;
+import org.apache.iceberg.spark.SparkWriteConf;
+import org.junit.Test;
+
+import java.lang.reflect.Field;
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.Assert.assertEquals;
+import static org.mockito.Mockito.CALLS_REAL_METHODS;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+public class IcebergWriteExecCapacityTest {
+  @Test
+  public void unitlessTablePropertiesDefaultToBytes() throws Exception {
+    for (String value : new String[] {"0", "1024", " 1024 "}) {
+      assertTableProperties(value, value.trim() + "B");
+    }
+  }
+
+  @Test
+  public void explicitTablePropertyUnitsArePreserved() throws Exception {
+    for (String value : new String[] {"1024B", "1KB", "1MB", "1mb", " 1 MB "}) 
{
+      assertTableProperties(value, value.trim());
+    }
+  }
+
+  @Test
+  public void absentTablePropertiesUseByteDefaults() throws Exception {
+    IcebergWriteExec exec = createExec(new HashMap<>(), 8192L);
+    assertEquals(
+        TableProperties.PARQUET_PAGE_SIZE_BYTES_DEFAULT + "B", 
exec.getParquetPageSizeBytes());
+    assertEquals(TableProperties.PARQUET_DICT_SIZE_BYTES_DEFAULT + "B", 
exec.getDictSizeBytes());
+  }
+
+  @Test
+  public void targetFileSizeUsesBytes() throws Exception {
+    for (long value :
+        new long[] {0L, 8192L, 
TableProperties.WRITE_TARGET_FILE_SIZE_BYTES_DEFAULT}) {
+      assertEquals(value + "B", createExec(new HashMap<>(), 
value).getTargetFileSizeBytes());
+    }
+  }
+
+  private void assertTableProperties(String value, String expected) throws 
Exception {
+    Map<String, String> properties = new HashMap<>();
+    properties.put(TableProperties.PARQUET_PAGE_SIZE_BYTES, value);
+    properties.put(TableProperties.PARQUET_DICT_SIZE_BYTES, value);
+    IcebergWriteExec exec = createExec(properties, 8192L);
+    assertEquals(expected, exec.getParquetPageSizeBytes());
+    assertEquals(expected, exec.getDictSizeBytes());
+  }
+
+  private IcebergWriteExec createExec(Map<String, String> properties, long 
targetFileSize)
+      throws Exception {
+    Table table = mock(Table.class);
+    when(table.properties()).thenReturn(properties);
+    SparkWriteConf writeConf = mock(SparkWriteConf.class);
+    when(writeConf.targetDataFileSize()).thenReturn(targetFileSize);
+    SparkWrite write = mock(SparkWrite.class);
+    setWriteField(write, "table", table);
+    setWriteField(write, "writeConf", writeConf);
+    IcebergWriteExec exec = mock(IcebergWriteExec.class, CALLS_REAL_METHODS);
+    when(exec.write()).thenReturn(write);
+    return exec;
+  }
+
+  private void setWriteField(SparkWrite write, String name, Object value) 
throws Exception {
+    Field field = SparkWrite.class.getDeclaredField(name);
+    field.setAccessible(true);
+    field.set(write, value);
+  }
+}
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 edeced4935..39f815aadd 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
@@ -18,7 +18,7 @@ package org.apache.gluten.execution.enhanced
 
 import org.apache.gluten.config.GlutenConfig.COLUMNAR_PARQUET_WRITE_BLOCK_SIZE
 import org.apache.gluten.config.GlutenIcebergConfig
-import org.apache.gluten.config.VeloxConfig.MAX_TARGET_FILE_SIZE_SESSION
+import org.apache.gluten.config.VeloxConfig.{MAX_TARGET_FILE_SIZE_SESSION, 
PARQUET_DICT_SIZE_BYTES, PARQUET_PAGE_SIZE_BYTES}
 import org.apache.gluten.execution._
 import org.apache.gluten.tags.EnhancedFeaturesTest
 
@@ -542,76 +542,92 @@ class VeloxIcebergSuite extends IcebergSuite {
       )
     }
   }
-  ignore("disabled test") {
-    test("iceberg native write respects target file size bytes") {
-      withTable("iceberg_small_target_tbl") {
-        spark.sql(
-          """
-            |CREATE TABLE iceberg_small_target_tbl (
-            |  id INT,
-            |  payload STRING
-            |) USING iceberg
-            |TBLPROPERTIES (
-            |  'write.format.default' = 'parquet',
-            |  'write.parquet.compression-codec' = 'uncompressed',
-            |  'write.parquet.row-group-size-bytes' = '4096',
-            |  'write.parquet.page-size-bytes' = '1024B',
-            |  'write.target-file-size-bytes' = '8192'
-            |)
-            |""".stripMargin)
-
-        checkAnswer(
-          spark.sql(
-            """
-              |SHOW TBLPROPERTIES iceberg_small_target_tbl
-              |('write.target-file-size-bytes')
-              |""".stripMargin),
-          Seq(Row("write.target-file-size-bytes", "8192"))
-        )
-
-        val df = spark.sql(
-          """
-            |INSERT INTO iceberg_small_target_tbl
-            |SELECT /*+ COALESCE(1) */
-            |  CAST(id AS INT),
-            |  concat(
-            |    CAST(id AS STRING),
-            |    '-',
-            |    sha2(CAST(id AS STRING), 256),
-            |    '-',
-            |    sha2(CAST(id + 1000 AS STRING), 256)
-            |  )
-            |FROM range(1000)
-            |""".stripMargin)
-
-        val commandPlan =
-          
df.queryExecution.executedPlan.asInstanceOf[CommandResultExec].commandPhysicalPlan
-
-        assert(commandPlan.isInstanceOf[VeloxIcebergAppendDataExec])
-
-        checkAnswer(
-          spark.sql("SELECT COUNT(*) FROM iceberg_small_target_tbl"),
-          Seq(Row(1000L)))
-
-        val files = spark.sql(
-          """
-            |SELECT file_size_in_bytes
-            |FROM default.iceberg_small_target_tbl.files
-            |""".stripMargin).collect().map(_.getLong(0))
-
-        assert(files.nonEmpty)
-
-        assert(
-          files.length > 1,
-          s"Expected write.target-file-size-bytes=8192 to create multiple 
files, " +
-            s"but got files=${files.mkString("[", ", ", "]")}")
-
-        assert(
-          files.max < 64L * 1024L,
-          s"Expected small target file size to keep max file size reasonably 
small, " +
-            s"but got files=${files.mkString("[", ", ", "]")}")
+  Seq(None, Some("8192"), Some("8KB"), Some("0")).foreach {
+    targetFileSizeOverride =>
+      val description = targetFileSizeOverride.map(v => s" with session 
override $v").getOrElse("")
+      test(s"iceberg native write respects target file size 
bytes$description") {
+        val overrides = targetFileSizeOverride.toSeq.flatMap {
+          value =>
+            Seq(
+              MAX_TARGET_FILE_SIZE_SESSION.key -> value,
+              PARQUET_PAGE_SIZE_BYTES.key -> "1024",
+              PARQUET_DICT_SIZE_BYTES.key -> "2048")
+        }
+        withSQLConf(overrides: _*) {
+          withTable("iceberg_small_target_tbl") {
+            spark.sql(
+              """
+                |CREATE TABLE iceberg_small_target_tbl (
+                |  id INT,
+                |  payload STRING
+                |) USING iceberg
+                |TBLPROPERTIES (
+                |  'write.format.default' = 'parquet',
+                |  'write.parquet.compression-codec' = 'uncompressed',
+                |  'write.parquet.page-size-bytes' = '1024',
+                |  'write.target-file-size-bytes' = '8192'
+                |)
+                |""".stripMargin)
+
+            checkAnswer(
+              spark.sql(
+                """
+                  |SHOW TBLPROPERTIES iceberg_small_target_tbl
+                  |('write.target-file-size-bytes')
+                  |""".stripMargin),
+              Seq(Row("write.target-file-size-bytes", "8192"))
+            )
+
+            val df = spark.sql(
+              """
+                |INSERT INTO iceberg_small_target_tbl
+                |SELECT /*+ COALESCE(1) */
+                |  CAST(id AS INT),
+                |  concat(
+                |    CAST(id AS STRING),
+                |    '-',
+                |    sha2(CAST(id AS STRING), 256),
+                |    '-',
+                |    sha2(CAST(id + 1000 AS STRING), 256)
+                |  )
+                |FROM range(0, 1000, 1, 16)
+                |""".stripMargin)
+
+            val commandPlan =
+              
df.queryExecution.executedPlan.asInstanceOf[CommandResultExec].commandPhysicalPlan
+
+            assert(commandPlan.isInstanceOf[VeloxIcebergAppendDataExec])
+
+            checkAnswer(
+              spark.sql("SELECT COUNT(*) FROM iceberg_small_target_tbl"),
+              Seq(Row(1000L)))
+
+            val files = spark.sql(
+              """
+                |SELECT file_size_in_bytes
+                |FROM default.iceberg_small_target_tbl.files
+                |""".stripMargin).collect().map(_.getLong(0))
+
+            assert(files.nonEmpty)
+
+            if (targetFileSizeOverride.contains("0")) {
+              assert(
+                files.length == 1,
+                s"Expected the session override to disable file rotation: 
${files.mkString(", ")}")
+            } else {
+              assert(
+                files.length > 1,
+                s"Expected write.target-file-size-bytes=8192 to create 
multiple files, " +
+                  s"but got files=${files.mkString("[", ", ", "]")}")
+
+              assert(
+                files.max < 64L * 1024L,
+                s"Expected small target file size to keep max file size 
reasonably small, " +
+                  s"but got files=${files.mkString("[", ", ", "]")}")
+            }
+          }
+        }
       }
-    }
   }
 
   test("iceberg parquet writer respects dictionary page size bytes") {
@@ -685,7 +701,7 @@ class VeloxIcebergSuite extends IcebergSuite {
                      |TBLPROPERTIES (
                      |  'write.format.default' = 'parquet',
                      |  'write.parquet.compression-codec' = 'uncompressed',
-                     |  'write.parquet.dict-size-bytes' = '1B'
+                     |  'write.parquet.dict-size-bytes' = '1'
                      |)
                      |""".stripMargin)
 
@@ -715,7 +731,7 @@ class VeloxIcebergSuite extends IcebergSuite {
         )
         assert(
           encodings.contains(Encoding.PLAIN),
-          s"Expected write.parquet.dict-size-bytes=1B to make later data pages 
fall back " +
+          s"Expected write.parquet.dict-size-bytes=1 to make later data pages 
fall back " +
             s"to PLAIN, but got encodings=${encodings.mkString("[", ", ", 
"]")}"
         )
       }
diff --git a/cpp/velox/utils/ConfigExtractor.cc 
b/cpp/velox/utils/ConfigExtractor.cc
index 77bf48fef9..0d551a93eb 100644
--- a/cpp/velox/utils/ConfigExtractor.cc
+++ b/cpp/velox/utils/ConfigExtractor.cc
@@ -265,7 +265,7 @@ std::shared_ptr<facebook::velox::config::ConfigBase> 
createHiveConnectorSessionC
   
configs[facebook::velox::connector::hive::HiveConfig::kReadTimestampUnitSession]
 = std::string("6");
   
configs[facebook::velox::connector::hive::HiveConfig::kMaxPartitionsPerWritersSession]
 =
       conf->get<std::string>(kMaxPartitions, "10000");
-  
configs[facebook::velox::connector::hive::HiveConfig::kParquetMaxTargetFileSize]
 =
+  
configs[facebook::velox::connector::hive::HiveConfig::kParquetMaxTargetFileSizeSession]
 =
       conf->get<std::string>(kParquetMaxTargetFileSize, "0B"); // 0 means no 
limit on target file size
   
configs[facebook::velox::connector::hive::HiveConfig::kIgnoreMissingFilesSession]
 =
       conf->get<bool>(kIgnoreMissingFiles, false) ? "true" : "false";
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 aa8671f109..284051c47f 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
@@ -53,9 +53,8 @@ trait IcebergWriteExec extends ColumnarV2TableWriteExec {
 
   protected def getParquetPageSizeBytes: String = {
     val tableProps = IcebergWriteUtil.getTable(write).properties()
-    tableProps.getOrDefault(
-      normalizeCapacityString(PARQUET_PAGE_SIZE_BYTES),
-      normalizeCapacityString(PARQUET_PAGE_SIZE_BYTES_DEFAULT.toString))
+    normalizeCapacityString(
+      tableProps.getOrDefault(PARQUET_PAGE_SIZE_BYTES, 
PARQUET_PAGE_SIZE_BYTES_DEFAULT.toString))
   }
 
   protected def getParquetPageRowLimit: String = {
@@ -64,7 +63,7 @@ trait IcebergWriteExec extends ColumnarV2TableWriteExec {
   }
 
   protected def getTargetFileSizeBytes: String = {
-    IcebergWriteUtil.getWriteConf(write).targetDataFileSize().toString
+    
normalizeCapacityString(IcebergWriteUtil.getWriteConf(write).targetDataFileSize().toString)
   }
 
   protected def getParquetRowGroupSizeBytes: String = {
@@ -77,9 +76,8 @@ trait IcebergWriteExec extends ColumnarV2TableWriteExec {
 
   protected def getDictSizeBytes: String = {
     val tableProps = IcebergWriteUtil.getTable(write).properties()
-    tableProps.getOrDefault(
-      normalizeCapacityString(PARQUET_DICT_SIZE_BYTES),
-      normalizeCapacityString(PARQUET_DICT_SIZE_BYTES_DEFAULT.toString))
+    normalizeCapacityString(
+      tableProps.getOrDefault(PARQUET_DICT_SIZE_BYTES, 
PARQUET_DICT_SIZE_BYTES_DEFAULT.toString))
   }
 
   protected def getPartitionSpec: PartitionSpec = {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to