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 09d90e2022 [VL] Honor Iceberg Parquet compression level and page 
version (#13073)
09d90e2022 is described below

commit 09d90e2022fb3e9f093f3a30971a90f80ce5ad11
Author: inf <[email protected]>
AuthorDate: Thu Oct 1 12:12:25 2026 +0000

    [VL] Honor Iceberg Parquet compression level and page version (#13073)
---
 .../execution/AbstractIcebergWriteExec.scala       |  22 ++++-
 .../execution/enhanced/VeloxIcebergSuite.scala     | 110 +++++++++++++++++++++
 cpp/velox/compute/VeloxBackend.cc                  |   4 +-
 cpp/velox/tests/iceberg/IcebergWriteTest.cc        |  35 ++++++-
 cpp/velox/utils/VeloxWriterUtils.cc                |  11 +++
 cpp/velox/utils/VeloxWriterUtils.h                 |   7 ++
 docs/get-started/VeloxIceberg.md                   |  12 ++-
 .../apache/gluten/execution/IcebergWriteExec.scala |  15 ++-
 8 files changed, 204 insertions(+), 12 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 6df3b105f5..2fe2ce06ff 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
@@ -36,6 +36,12 @@ abstract class AbstractIcebergWriteExec extends 
IcebergWriteExec {
   private val parquetPageRowLimitSession =
     "spark.gluten.sql.columnar.backend.velox.parquet_writer_page_row_limit"
 
+  private val parquetPageVersionSession =
+    "spark.gluten.sql.columnar.backend.velox.parquet_writer_datapage_version"
+
+  private val parquetCompressionLevelSession =
+    "spark.gluten.sql.columnar.backend.velox.parquet_writer_compression_level"
+
   // the writer factory works for both batch and streaming
   private def createIcebergDataWriteFactory(schema: StructType): 
IcebergDataWriteFactory = {
     val writeSchema = IcebergWriteUtil.getWriteSchema(write)
@@ -53,7 +59,6 @@ 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 {
@@ -61,11 +66,24 @@ abstract class AbstractIcebergWriteExec extends 
IcebergWriteExec {
         val overrideValue = SQLConf.get.getConfString(key, null)
         if (overrideValue == null) {
           icebergProperties.put(key, value)
-        } else if (key != parquetPageRowLimitSession) {
+        } else {
           icebergProperties.put(key, normalizeCapacityString(overrideValue))
         }
     }
 
+    Seq(
+      parquetPageRowLimitSession -> getParquetPageRowLimit,
+      parquetPageVersionSession -> getParquetPageVersion).foreach {
+      case (key, value) =>
+        icebergProperties.put(key, SQLConf.get.getConfString(key, value))
+    }
+
+    if (Seq("gzip", "zstd").exists(_.equalsIgnoreCase(getCodec))) {
+      Option(SQLConf.get.getConfString(parquetCompressionLevelSession, null))
+        .orElse(getParquetCompressionLevel)
+        .foreach(level => 
icebergProperties.put(parquetCompressionLevelSession, level))
+    }
+
     IcebergDataWriteFactory(
       filteredSchema,
       getFileFormat(IcebergWriteUtil.getFileFormat(write)),
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 39f815aadd..145e9f6a49 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
@@ -738,6 +738,116 @@ class VeloxIcebergSuite extends IcebergSuite {
     }
   }
 
+  Seq(
+    ("v1", None),
+    ("v2", None),
+    ("V2", None),
+    ("v1", Some("V2")),
+    ("v2", Some("V1"))).foreach {
+    case (version, sessionVersion) =>
+      test(s"iceberg parquet page version $version with session override 
$sessionVersion") {
+        val overrides = sessionVersion.map {
+          
"spark.gluten.sql.columnar.backend.velox.parquet_writer_datapage_version" -> _
+        }.toSeq
+        withSQLConf(overrides: _*) {
+          withTable("iceberg_page_version") {
+            spark.sql(s"""
+          CREATE TABLE iceberg_page_version (id BIGINT) USING iceberg
+          TBLPROPERTIES ('write.parquet.page-version' = '$version')
+        """)
+            val df = spark.sql("INSERT INTO iceberg_page_version SELECT id 
FROM range(1000)")
+            assert(
+              df.queryExecution.executedPlan
+                .asInstanceOf[CommandResultExec]
+                .commandPhysicalPlan
+                .isInstanceOf[VeloxIcebergAppendDataExec])
+            val files =
+              spark.sql("SELECT file_path FROM 
default.iceberg_page_version.files").collect()
+            assert(files.nonEmpty)
+            files.foreach {
+              file =>
+                val reader = ParquetFileReader.open(HadoopInputFile.fromPath(
+                  new Path(file.getString(0)),
+                  spark.sparkContext.hadoopConfiguration))
+                try {
+                  val column = 
reader.getFooter.getFileMetaData.getSchema.getColumns.get(0)
+                  val page = 
reader.readNextRowGroup().getPageReader(column).readPage()
+                  assert(page != null)
+                  assert(
+                    page.isInstanceOf[DataPageV2] ==
+                      sessionVersion.getOrElse(version).equalsIgnoreCase("v2"))
+                } finally {
+                  reader.close()
+                }
+            }
+            checkAnswer(spark.sql("SELECT count(*) FROM 
iceberg_page_version"), Seq(Row(1000L)))
+          }
+        }
+      }
+  }
+
+  Seq("gzip", "zstd").foreach {
+    codec =>
+      test(s"iceberg parquet $codec compression level and write option 
precedence") {
+        withSQLConf("spark.sql.shuffle.partitions" -> "1") {
+          withTable("iceberg_compression_level") {
+            spark.sql(s"""
+            CREATE TABLE iceberg_compression_level (value STRING) USING iceberg
+            TBLPROPERTIES (
+              'write.parquet.compression-codec' = '$codec',
+              'write.parquet.compression-level' = '1',
+              'write.parquet.dict-size-bytes' = '1B')
+          """)
+            val data = spark.range(0, 10000, 1, 1)
+              .selectExpr("concat(cast(id as string), repeat('abcdefghij', 
100)) AS value")
+
+            def writeSize(level: Option[String]): Long = {
+              spark.sql("TRUNCATE TABLE iceberg_compression_level")
+              val writer = data.writeTo("iceberg_compression_level")
+              level.foreach(writer.option("compression-level", _))
+              
TestUtils.checkExecutedPlanContains[VeloxIcebergAppendDataExec](spark) {
+                writer.append()
+              }
+              checkAnswer(
+                spark.sql("SELECT count(*) FROM iceberg_compression_level"),
+                Seq(Row(10000L)))
+              spark.sql(
+                "SELECT sum(file_size_in_bytes) FROM 
default.iceberg_compression_level.files")
+                .head().getLong(0)
+            }
+
+            val low = writeSize(None)
+            val high = writeSize(Some("9"))
+            assert(high < low, s"Expected level 9 to compress better than 
level 1: $high >= $low")
+            spark.sql("""
+            ALTER TABLE iceberg_compression_level SET TBLPROPERTIES
+              ('write.parquet.compression-level' = '9')
+          """)
+            assert(writeSize(None) == high)
+            withSQLConf("spark.sql.iceberg.compression-level" -> "1") {
+              assert(writeSize(None) == low)
+              assert(writeSize(Some("9")) == high)
+            }
+            val glutenCompressionLevel =
+              
"spark.gluten.sql.columnar.backend.velox.parquet_writer_compression_level"
+            withSQLConf(
+              glutenCompressionLevel -> "1",
+              "spark.sql.iceberg.compression-level" -> "9") {
+              assert(writeSize(None) == low)
+              assert(writeSize(Some("9")) == low)
+            }
+            spark.sql("""
+            ALTER TABLE iceberg_compression_level UNSET TBLPROPERTIES
+              ('write.parquet.compression-level')
+          """)
+            withSQLConf(glutenCompressionLevel -> "1") {
+              assert(writeSize(None) == low)
+            }
+          }
+        }
+      }
+  }
+
   test("iceberg table page row limit") {
     val table = "iceberg_page_row_limit"
 
diff --git a/cpp/velox/compute/VeloxBackend.cc 
b/cpp/velox/compute/VeloxBackend.cc
index 074aec9df4..a3e1dbd9cc 100644
--- a/cpp/velox/compute/VeloxBackend.cc
+++ b/cpp/velox/compute/VeloxBackend.cc
@@ -51,6 +51,7 @@
 #include "shuffle/ArrowShuffleDictionaryWriter.h"
 #include "udf/UdfLoader.h"
 #include "utils/Exception.h"
+#include "utils/VeloxWriterUtils.h"
 #include "velox/common/caching/SsdCache.h"
 #include "velox/common/file/FileSystems.h"
 #include "velox/connectors/hive/BufferedInputBuilder.h"
@@ -63,7 +64,6 @@
 #include 
"velox/connectors/hive/storage_adapters/hdfs/RegisterHdfsFileSystem.h" // 
@manual
 #include "velox/dwio/orc/reader/OrcReader.h"
 #include "velox/dwio/parquet/RegisterParquetReader.h"
-#include "velox/dwio/parquet/RegisterParquetWriter.h"
 #include "velox/serializers/PrestoSerializer.h"
 
 DECLARE_bool(velox_exception_user_stacktrace_enabled);
@@ -258,7 +258,7 @@ void VeloxBackend::init(
 
   velox::dwio::common::registerFileSinks();
   velox::parquet::registerParquetReaderFactory();
-  velox::parquet::registerParquetWriterFactory();
+  
velox::dwio::common::registerWriterFactory(std::make_shared<GlutenParquetWriterFactory>());
   velox::orc::registerOrcReaderFactory();
   
velox::exec::ExprToSubfieldFilterParser::registerParser(std::make_unique<SparkExprToSubfieldFilterParser>(
       backendConf_->get<bool>(kScanBloomFilterPushdownEnabled, 
kScanBloomFilterPushdownEnabledDefault)));
diff --git a/cpp/velox/tests/iceberg/IcebergWriteTest.cc 
b/cpp/velox/tests/iceberg/IcebergWriteTest.cc
index 740fa66f3f..b3a97fc9d3 100644
--- a/cpp/velox/tests/iceberg/IcebergWriteTest.cc
+++ b/cpp/velox/tests/iceberg/IcebergWriteTest.cc
@@ -17,7 +17,10 @@
 
 #include "compute/iceberg/IcebergWriter.h"
 #include "memory/VeloxColumnarBatch.h"
-#include "velox/dwio/parquet/RegisterParquetWriter.h"
+#include "utils/ConfigExtractor.h"
+#include "utils/VeloxWriterUtils.h"
+#include "velox/connectors/hive/FileConnectorUtil.h"
+#include "velox/connectors/hive/HiveConfig.h"
 #include "velox/exec/tests/utils/TempDirectoryPath.h"
 #include "velox/vector/tests/utils/VectorTestBase.h"
 
@@ -30,7 +33,7 @@ class VeloxIcebergWriteTest : public ::testing::Test, public 
test::VectorTestBas
  protected:
   static void SetUpTestCase() {
     
memory::MemoryManager::testingSetInstance(memory::MemoryManager::Options{});
-    parquet::registerParquetWriterFactory();
+    
dwio::common::registerWriterFactory(std::make_shared<GlutenParquetWriterFactory>());
     Type::registerSerDe();
     dwio::common::registerFileSinks();
     filesystems::registerLocalFileSystem();
@@ -40,6 +43,32 @@ class VeloxIcebergWriteTest : public ::testing::Test, public 
test::VectorTestBas
   std::shared_ptr<memory::MemoryPool> connectorPool_ = 
rootPool_->addAggregateChild("connector");
 };
 
+TEST_F(VeloxIcebergWriteTest, parquetWriterOptions) {
+  GlutenParquetWriterFactory factory;
+  const config::ConfigBase empty(std::unordered_map<std::string, 
std::string>{});
+  auto defaults = 
std::static_pointer_cast<parquet::ParquetWriterOptions>(factory.createFormatOptions(empty,
 empty));
+  EXPECT_EQ(defaults->codecOptions, nullptr);
+  EXPECT_FALSE(defaults->useParquetDataPageV2.value_or(false));
+
+  const connector::hive::HiveConfig hiveConfig(
+      std::make_shared<config::ConfigBase>(std::unordered_map<std::string, 
std::string>{}));
+  for (auto level : {-5, 1, 9}) {
+    for (const auto& version : {"V1", "V2"}) {
+      auto sparkConf = 
std::make_shared<config::ConfigBase>(std::unordered_map<std::string, 
std::string>{
+          
{"spark.gluten.sql.columnar.backend.velox.parquet_writer_compression_level", 
std::to_string(level)},
+          
{"spark.gluten.sql.columnar.backend.velox.parquet_writer_datapage_version", 
version}});
+      auto session = createHiveConnectorSessionConfig(sparkConf);
+      auto scopedConfigs =
+          connector::hive::makeFormatScopedConfigs(hiveConfig, *session, 
dwio::common::FileFormat::PARQUET);
+      auto options = std::static_pointer_cast<parquet::ParquetWriterOptions>(
+          factory.createFormatOptions(scopedConfigs.connectorConfig, 
scopedConfigs.sessionProperties));
+      ASSERT_NE(options->codecOptions, nullptr);
+      EXPECT_EQ(options->codecOptions->compressionLevel, level);
+      EXPECT_EQ(options->useParquetDataPageV2.value(), std::string(version) == 
"V2");
+    }
+  }
+}
+
 TEST_F(VeloxIcebergWriteTest, write) {
   auto vector = makeRowVector({makeFlatVector<int8_t>({1, 2}), 
makeFlatVector<int16_t>({1, 2})});
   auto tmpPath = tmpDir_->getPath();
@@ -64,7 +93,7 @@ TEST_F(VeloxIcebergWriteTest, write) {
       partitionSpec,
       root,
       std::unordered_map<std::string, std::string>(),
-      pool_,
+      rootPool_,
       connectorPool_);
   auto batch = VeloxColumnarBatch(vector);
   writer->write(batch);
diff --git a/cpp/velox/utils/VeloxWriterUtils.cc 
b/cpp/velox/utils/VeloxWriterUtils.cc
index 68987edd99..a2491836ea 100644
--- a/cpp/velox/utils/VeloxWriterUtils.cc
+++ b/cpp/velox/utils/VeloxWriterUtils.cc
@@ -38,6 +38,17 @@ const int32_t kGzipWindowBits4k = 12;
 const int32_t kZSTDDefaultCompressionLevel = 3;
 } // namespace
 
+std::shared_ptr<dwio::common::FormatSpecificOptions> 
GlutenParquetWriterFactory::createFormatOptions(
+    const config::ConfigBase& connectorConfig,
+    const config::ConfigBase& session) const {
+  auto options = ParquetWriterFactory::createFormatOptions(connectorConfig, 
session);
+  if (auto level = session.get<int32_t>("writer_compression_level")) {
+    auto parquetOptions = 
std::static_pointer_cast<ParquetWriterOptions>(options);
+    parquetOptions->codecOptions = 
std::make_shared<parquet::arrow::util::CodecOptions>(*level);
+  }
+  return options;
+}
+
 std::shared_ptr<facebook::velox::dwio::common::WriterOptions> 
makeParquetWriteOption(
     const std::unordered_map<std::string, std::string>& sparkConfs) {
   int64_t maxRowGroupBytes = 134217728; // 128MB
diff --git a/cpp/velox/utils/VeloxWriterUtils.h 
b/cpp/velox/utils/VeloxWriterUtils.h
index 56bd5b2e27..50622776e5 100644
--- a/cpp/velox/utils/VeloxWriterUtils.h
+++ b/cpp/velox/utils/VeloxWriterUtils.h
@@ -23,6 +23,13 @@
 
 namespace gluten {
 
+class GlutenParquetWriterFactory : public 
facebook::velox::parquet::ParquetWriterFactory {
+ public:
+  std::shared_ptr<facebook::velox::dwio::common::FormatSpecificOptions> 
createFormatOptions(
+      const facebook::velox::config::ConfigBase& connectorConfig,
+      const facebook::velox::config::ConfigBase& session) const override;
+};
+
 std::shared_ptr<facebook::velox::dwio::common::WriterOptions> 
makeParquetWriteOption(
     const std::unordered_map<std::string, std::string>& sparkConfs);
 
diff --git a/docs/get-started/VeloxIceberg.md b/docs/get-started/VeloxIceberg.md
index 21cff1df5f..6b81da3e77 100644
--- a/docs/get-started/VeloxIceberg.md
+++ b/docs/get-started/VeloxIceberg.md
@@ -104,7 +104,10 @@ the added column name is same to the deleted column, the 
scan will fall back.
 | spark.gluten.sql.columnar.iceberg.enableNativeRead | true | Enable 
offloading Iceberg scans to the native backend. When disabled, Iceberg scans 
fall back to vanilla Spark while scans of other formats stay offloaded. |
 | spark.gluten.sql.columnar.iceberg.enableNativeWrite | true | Enable 
offloading Iceberg writes to the native backend. When disabled, Iceberg writes 
fall back to vanilla Spark. Note the Velox backend additionally requires 
`spark.gluten.sql.enable.enhancedFeatures` to be enabled. |
 | spark.gluten.sql.columnar.parquet.write.blockSize | 128MB | Target Parquet 
row-group size for native writes. When explicitly set, overrides the Iceberg 
table property `write.parquet.row-group-size-bytes`. Accepts byte counts or 
sizes such as `32MB` and `1GB`. |
-Both options are runtime modifiable, so they can be flipped per session with 
`SET`.
+| spark.gluten.sql.columnar.backend.velox.parquet_writer_compression_level | 
(unset) | Overrides Iceberg's Parquet compression level for gzip and zstd. |
+| spark.gluten.sql.columnar.backend.velox.parquet_writer_datapage_version | 
(unset) | Overrides `write.parquet.page-version`; accepts `V1` or `V2`. |
+
+These options can be changed per session with `SET`.
 
 ### Catalogs
 All the catalog configurations are transparent to Gluten
@@ -136,7 +139,7 @@ The "Gluten Support" column is now ready to be populated 
with:
 | spark.wap.id | null | Write-Audit-Publish snapshot staging ID | |
 | spark.wap.branch | null | WAP branch name for snapshot commit | |
 | spark.sql.iceberg.compression-codec | Table default | Write compression 
codec (e.g., zstd, snappy) | |
-| spark.sql.iceberg.compression-level | Table default | Compression level for 
Parquet/Avro | |
+| spark.sql.iceberg.compression-level | Table default | Compression level for 
Parquet/Avro |⚠️ Parquet (gzip, zstd) only|
 | spark.sql.iceberg.compression-strategy | Table default | Compression 
strategy for ORC | |
 | spark.sql.iceberg.data-planning-mode | AUTO | Scan planning mode for data 
files (AUTO, LOCAL, DISTRIBUTED) | |
 | spark.sql.iceberg.delete-planning-mode | AUTO | Scan planning mode for 
delete files (AUTO, LOCAL, DISTRIBUTED) | |
@@ -177,7 +180,7 @@ The "Gluten Support" column is now ready to be populated 
with:
 | isolation-level | null | Desired isolation level for Dataframe overwrite 
operations. null => no checks (for idempotent writes), serializable => check 
for concurrent inserts or deletes in destination partitions, snapshot => checks 
for concurrent deletes in destination partitions. | |
 | validate-from-snapshot-id | null | If isolation level is set, id of base 
snapshot from which to check concurrent write conflicts into a table. Should be 
the snapshot before any reads from the table. Can be obtained via Table API or 
Snapshots table. If null, the table's oldest known snapshot is used. | |
 | compression-codec | Table write.(fileformat).compression-codec | Overrides 
this table's compression codec for this write | |
-| compression-level | Table write.(fileformat).compression-level | Overrides 
this table's compression level for Parquet and Avro tables for this write | |
+| compression-level | Table write.(fileformat).compression-level | Overrides 
the Iceberg session and table compression levels for this write |⚠️ Parquet 
(gzip, zstd) only|
 | compression-strategy | Table write.orc.compression-strategy | Overrides this 
table's compression strategy for ORC tables for this write | |
 | distribution-mode | See Spark Writes for defaults | Override this table's 
distribution mode for this write |🚫|
 | delete-granularity | file | Override this table's delete granularity for 
this write | |
@@ -206,10 +209,11 @@ extracted from 
https://iceberg.apache.org/docs/latest/configuration/
 | write.delete.format.default | data file format | Default delete file format 
for the table; parquet, avro, or orc |  |
 | write.parquet.row-group-size-bytes | 134217728 (128 MB) | Target Parquet 
row-group size in bytes. Overridden by an explicitly set 
`spark.gluten.sql.columnar.parquet.write.blockSize`. |✅|
 | write.parquet.page-size-bytes | 1048576 (1 MB) | Parquet page size |✅|
+| write.parquet.page-version | v1 | Parquet data page version: `v1` or `v2` 
(case-insensitive) |✅|
 | write.parquet.page-row-limit | 20000 | Parquet page row limit |  |
 | write.parquet.dict-size-bytes | 2097152 (2 MB) | Parquet dictionary page 
size |  |
 | write.parquet.compression-codec | zstd | Parquet compression codec: zstd, 
lz4, gzip, snappy, uncompressed. **Note:** Native writes fall back to Spark for 
brotli, lzo, lz4raw, and lz4_raw |⚠️|
-| write.parquet.compression-level | null | Parquet compression level |  |
+| write.parquet.compression-level | null | Parquet compression level; unset 
uses the codec default. Overridden by `spark.sql.iceberg.compression-level` and 
the per-write `compression-level` option. |✅ gzip, zstd|
 | write.parquet.bloom-filter-enabled.column.col1 | (not set) | Hint to parquet 
to write a bloom filter for the column: 'col1' |  |
 | write.parquet.bloom-filter-max-bytes | 1048576 (1 MB) | The maximum number 
of bytes for a bloom filter bitset |  |
 | write.parquet.bloom-filter-fpp.column.col1 | 0.01 | The false positive 
probability for a bloom filter applied to 'col1' (must > 0.0 and < 1.0) |  |
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 284051c47f..cf5c047514 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_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.TableProperties._
 import org.apache.iceberg.avro.AvroSchemaUtil
 import org.apache.iceberg.spark.source.IcebergWriteUtil
 import org.apache.iceberg.types.Type.TypeID
@@ -51,6 +51,19 @@ trait IcebergWriteExec extends ColumnarV2TableWriteExec {
     } else codec.toLowerCase(Locale.ROOT)
   }
 
+  protected def getParquetCompressionLevel: Option[String] = {
+    
Option(IcebergWriteUtil.getWriteProperty(write).get(PARQUET_COMPRESSION_LEVEL))
+      
.orElse(Option(IcebergWriteUtil.getTable(write).properties().get(PARQUET_COMPRESSION_LEVEL)))
+  }
+
+  protected def getParquetPageVersion: String = {
+    IcebergWriteUtil
+      .getTable(write)
+      .properties()
+      .getOrDefault("write.parquet.page-version", "v1")
+      .toUpperCase(Locale.ROOT)
+  }
+
   protected def getParquetPageSizeBytes: String = {
     val tableProps = IcebergWriteUtil.getTable(write).properties()
     normalizeCapacityString(


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

Reply via email to