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 16d1b87ab8 [VL][Iceberg] Fix iceberg codec validation (#12364)
16d1b87ab8 is described below
commit 16d1b87ab897b1a1279c3fd9abd0279bdd03977c
Author: inf <[email protected]>
AuthorDate: Tue Sep 22 16:38:52 2026 +0000
[VL][Iceberg] Fix iceberg codec validation (#12364)
---
.../execution/enhanced/VeloxIcebergSuite.scala | 68 +++++++++++++++++++++-
docs/get-started/VeloxIceberg.md | 2 +-
.../apache/gluten/execution/IcebergWriteExec.scala | 9 ++-
3 files changed, 74 insertions(+), 5 deletions(-)
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 3ad8c7ffbc..edeced4935 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
@@ -23,7 +23,7 @@ import org.apache.gluten.execution._
import org.apache.gluten.tags.EnhancedFeaturesTest
import org.apache.spark.sql.{DataFrame, Row}
-import org.apache.spark.sql.execution.CommandResultExec
+import org.apache.spark.sql.execution.{CommandExecutionMode, CommandResultExec}
import org.apache.spark.sql.execution.GlutenImplicits._
import org.apache.spark.sql.execution.datasources.v2.AppendDataExec
import org.apache.spark.sql.execution.streaming.MemoryStream
@@ -36,6 +36,8 @@ import
org.apache.iceberg.shaded.org.apache.parquet.column.page.{DataPage, DataP
import org.apache.iceberg.shaded.org.apache.parquet.hadoop.ParquetFileReader
import org.apache.iceberg.shaded.org.apache.parquet.hadoop.util.HadoopInputFile
+import java.util.Locale
+
import scala.jdk.CollectionConverters._
@EnhancedFeaturesTest
@@ -43,6 +45,70 @@ class VeloxIcebergSuite extends IcebergSuite {
import testImplicits._
+ test("iceberg write falls back for unsupported compression codecs") {
+ val codecs = Seq("brotli", "lzo", "lz4raw", "lz4_raw")
+ (codecs ++ codecs.map(_.toUpperCase(Locale.ROOT))).foreach {
+ codec =>
+ withTable("iceberg_codec_test") {
+ spark.sql(s"""
+ |CREATE TABLE iceberg_codec_test (id INT, data STRING)
USING iceberg
+ |TBLPROPERTIES ('write.parquet.compression-codec' =
'$codec')
+ |""".stripMargin)
+ // Plan without executing: fallback codecs may require optional
Hadoop libraries.
+ val logicalPlan = spark.sessionState.sqlParser.parsePlan(
+ "INSERT INTO iceberg_codec_test VALUES (1, 'test')")
+ val plan = spark.sessionState
+ .executePlan(logicalPlan, CommandExecutionMode.SKIP)
+ .executedPlan
+ val append = plan.collectFirst { case a: AppendDataExec => a }
+ assert(append.isDefined, s"Expected fallback for codec $codec:
$plan")
+ assert(!plan.exists(_.isInstanceOf[VeloxIcebergAppendDataExec]))
+ val validation =
VeloxIcebergAppendDataExec(append.get).doValidateInternal()
+ assert(!validation.ok())
+ assert(validation.reason().contains("Codec unsupported"),
validation.reason())
+ }
+ }
+ }
+
+ test("iceberg write uses supported compression codecs") {
+ val codecs = Seq("snappy", "gzip", "zstd", "lz4", "uncompressed")
+ (codecs ++ codecs.map(_.toUpperCase(Locale.ROOT))).foreach {
+ codec =>
+ withTable("iceberg_codec_test") {
+ spark.sql(s"""
+ |CREATE TABLE iceberg_codec_test (id INT, data STRING)
USING iceberg
+ |TBLPROPERTIES ('write.parquet.compression-codec' =
'$codec')
+ |""".stripMargin)
+ TestUtils.checkExecutedPlanContains[VeloxIcebergAppendDataExec](
+ spark,
+ "INSERT INTO iceberg_codec_test VALUES (1, 'test')")
+ checkAnswer(spark.sql("SELECT * FROM iceberg_codec_test"),
Seq(Row(1, "test")))
+ val files = spark.sql("SELECT file_path FROM
default.iceberg_codec_test.files").collect()
+ assert(files.nonEmpty)
+ val expectedCodec = codec.toUpperCase(Locale.ROOT) match {
+ case "LZ4" => "LZ4_RAW"
+ case other => other
+ }
+ files.foreach {
+ file =>
+ val input = HadoopInputFile.fromPath(
+ new Path(file.getString(0)),
+ spark.sessionState.newHadoopConf())
+ val reader = ParquetFileReader.open(input)
+ try {
+ val columns =
reader.getFooter.getBlocks.asScala.flatMap(_.getColumns.asScala)
+ assert(columns.nonEmpty)
+ assert(
+ columns.forall(_.getCodec.name() == expectedCodec),
+ s"Expected $expectedCodec compression for codec $codec")
+ } finally {
+ reader.close()
+ }
+ }
+ }
+ }
+ }
+
test("iceberg insert") {
withTable("iceberg_tb2") {
spark.sql("""
diff --git a/docs/get-started/VeloxIceberg.md b/docs/get-started/VeloxIceberg.md
index d8454e8998..21cff1df5f 100644
--- a/docs/get-started/VeloxIceberg.md
+++ b/docs/get-started/VeloxIceberg.md
@@ -208,7 +208,7 @@ extracted from
https://iceberg.apache.org/docs/latest/configuration/
| write.parquet.page-size-bytes | 1048576 (1 MB) | Parquet page size |✅|
| 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,
brotli, lz4, gzip, snappy, uncompressed | |
+| 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.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 | |
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 99638c064c..aa8671f109 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
@@ -24,6 +24,8 @@ import org.apache.iceberg.avro.AvroSchemaUtil
import org.apache.iceberg.spark.source.IcebergWriteUtil
import org.apache.iceberg.types.Type.TypeID
+import java.util.Locale
+
import scala.collection.JavaConverters._
trait IcebergWriteExec extends ColumnarV2TableWriteExec {
@@ -46,7 +48,7 @@ trait IcebergWriteExec extends ColumnarV2TableWriteExec {
}
if (codec.equalsIgnoreCase("uncompressed")) {
"none"
- } else codec
+ } else codec.toLowerCase(Locale.ROOT)
}
protected def getParquetPageSizeBytes: String = {
@@ -126,8 +128,9 @@ trait IcebergWriteExec extends ColumnarV2TableWriteExec {
}
val codec = getCodec
- if (Seq("brotli, lzo").contains(codec)) {
- return ValidationResult.failed("Not support this codec " + codec)
+ val unsupported = Set("brotli", "lzo", "lz4raw", "lz4_raw")
+ if (unsupported.contains(codec.toLowerCase(Locale.ROOT))) {
+ return ValidationResult.failed("Codec unsupported: " + codec)
}
if (query.output.exists(a =>
!AvroSchemaUtil.makeCompatibleName(a.name).equals(a.name))) {
return ValidationResult.failed("Not support the compatible column name")
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]