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]