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]