This is an automated email from the ASF dual-hosted git repository.
danny0405 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new 6b5ea33228bb feat(storage): default parquet codec to zstd for flink
and spark 3.5+ (#19685)
6b5ea33228bb is described below
commit 6b5ea33228bb52efe3dc373da7524e06180222b7
Author: Shuo Cheng <[email protected]>
AuthorDate: Fri Aug 28 22:47:25 2026 +0800
feat(storage): default parquet codec to zstd for flink and spark 3.5+
(#19685)
* feat(storage): default parquet codec to zstd for flink and spark 3.5+
---
LICENSE | 32 +++++++++++++++++
README.md | 6 ++++
.../org/apache/hudi/config/HoodieWriteConfig.java | 42 ++++++++++++++++++++++
.../org/apache/hudi/io/HoodieBinaryCopyHandle.java | 3 +-
.../hudi/metadata/HoodieMetadataWriteUtils.java | 1 +
.../apache/hudi/config/TestHoodieWriteConfig.java | 36 +++++++++++++++++++
.../metadata/TestHoodieMetadataWriteUtils.java | 14 ++++++++
.../row/HoodieRowDataFileWriterFactory.java | 2 +-
.../TestHoodieRowDataParquetConfigInjector.java | 2 ++
.../io/storage/HoodieSparkFileWriterFactory.java | 4 +--
.../hudi/common/config/HoodieStorageConfig.java | 10 ++++--
.../common/config/TestHoodieStorageConfig.java | 16 +++++++++
.../apache/hudi/utils/TestFlinkWriteClients.java | 11 ++++++
.../hadoop/HoodieAvroFileWriterFactory.java | 2 +-
.../common/testutils/HoodieCommonTestHarness.java | 3 +-
.../testutils/reader/HoodieFileSliceTestUtils.java | 3 +-
...oodieAvroFileWriterFactoryVariantInference.java | 1 +
.../TestHoodieAvroParquetConfigInjector.java | 9 +++++
.../hudi/metadata/TestHoodieTableMetadataUtil.java | 9 ++++-
.../hudi/hadoop/testutils/InputFormatTestUtil.java | 3 +-
hudi-spark-datasource/README.md | 10 ++++++
.../java/org/apache/hudi/TestDataSourceUtils.java | 26 ++++++++++++++
.../bundle-validation/spark_hadoop_mr/write.scala | 4 +++
packaging/hudi-flink-bundle/pom.xml | 4 ++-
24 files changed, 236 insertions(+), 17 deletions(-)
diff --git a/LICENSE b/LICENSE
index f0d0759bd101..90fb0dc0d919 100644
--- a/LICENSE
+++ b/LICENSE
@@ -287,6 +287,38 @@ SOFTWARE.
-------------------------------------------------------------------------------
+This product bundles zstd-jni (https://github.com/luben/zstd-jni), which is
+available under the BSD 2-Clause License:
+
+zstd-jni: JNI bindings to Zstd Library
+
+Copyright (c) 2015-present, Luben Karavelov/ All rights reserved.
+
+BSD License
+
+Redistribution and use in source and binary forms, with or without
modification,
+are permitted provided that the following conditions are met:
+
+* Redistributions of source code must retain the above copyright notice, this
+ list of conditions and the following disclaimer.
+
+* Redistributions in binary form must reproduce the above copyright notice,
this
+ list of conditions and the following disclaimer in the documentation and/or
+ other materials provided with the distribution.
+
+THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" AND
+ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED
+WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
+DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE
FOR
+ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES
+(INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES;
+LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON
+ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT
+(INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
+SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
+
+-------------------------------------------------------------------------------
+
This product includes code from Apache Hadoop
* org.apache.hudi.common.bloom.InternalDynamicBloomFilter.java adapted from
org.apache.hadoop.util.bloom.DynamicBloomFilter.java
diff --git a/README.md b/README.md
index f2ab0ca3caa3..ee41fcaa7704 100644
--- a/README.md
+++ b/README.md
@@ -139,6 +139,12 @@ Refer to the table below for building with different Spark
and Scala versions.
| `-Dspark4.2` | hudi-spark4.2-bundle_2.13 |
For Spark 4.2 and Scala 2.13 (Needs java 17) |
| `-Dspark3` | hudi-spark3-bundle_2.12 (legacy bundle name) |
For Spark 3.5.x and Scala 2.12 |
+Hudi uses ZSTD as the default Parquet compression codec with Spark 3.5 and
newer. Spark 3.3 and 3.4
+retain GZIP because Hudi's non-vectorized file-group reader uses parquet-java
1.12.x and can leak
+off-heap memory when reading ZSTD files
([PARQUET-2160](https://issues.apache.org/jira/browse/PARQUET-2160)).
+Upgrade to Spark 3.5 or newer before using ZSTD. Keeping GZIP as the write
default does not remove the
+risk when an older Spark runtime reads ZSTD files produced by another engine.
+
Please note that only Spark-related bundles, i.e., `hudi-spark-bundle`,
`hudi-utilities-bundle`,
`hudi-utilities-slim-bundle`, can be built using `scala-2.13` profile. Hudi
Flink bundle cannot be built
using `scala-2.13` profile. To build these bundles on Scala 2.13, use the
following command:
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
index 46a09baebc71..c12acdbe2964 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
@@ -134,6 +134,10 @@ import static
org.apache.hudi.table.marker.ConflictDetectionUtils.getDefaultEarl
description = "Configurations that control write behavior on Hudi tables.
These can be directly passed down from even "
+ "higher level frameworks (e.g Spark datasources, Flink sink) and
utilities (e.g Hudi Streamer).")
public class HoodieWriteConfig extends HoodieConfig {
+
+ private static final String GZIP_COMPRESSION_CODEC = "gzip";
+ private static final String ZSTD_COMPRESSION_CODEC = "zstd";
+ private static final String MIN_SPARK_VERSION_WITH_ZSTD_DEFAULT = "3.5.0";
private static final long serialVersionUID = 0L;
// This is a constant as is should never be changed via config (will
invalidate previous commits)
@@ -1429,6 +1433,40 @@ public class HoodieWriteConfig extends HoodieConfig {
return new Builder();
}
+ @VisibleForTesting
+ static String getDefaultParquetCompressionCodec(EngineType engineType) {
+ switch (engineType) {
+ case FLINK:
+ return ZSTD_COMPRESSION_CODEC;
+ case SPARK:
+ // Spark 3.5 and newer use ZSTD. Spark 3.3 and 3.4 retain GZIP because
their
+ // non-vectorized file-group reader uses parquet-java 1.12.x and can
leak off-heap
+ // memory when reading ZSTD files:
https://issues.apache.org/jira/browse/PARQUET-2160.
+ Option<String> sparkVersion = getSparkRuntimeVersion();
+ return sparkVersion.isPresent()
+ && StringUtils.compareVersions(sparkVersion.get(),
MIN_SPARK_VERSION_WITH_ZSTD_DEFAULT) >= 0
+ ? ZSTD_COMPRESSION_CODEC : GZIP_COMPRESSION_CODEC;
+ default:
+ // The Java client does not own its Parquet runtime: Parquet
dependencies are provided by
+ // the embedding application, and older Parquet versions use Hadoop
native ZSTD rather than
+ // zstd-jni. For example, the recommended Kafka HDFS Connector 10.1.0
uses Parquet 1.11.1,
+ // which risks leaking memory when reading ZSTD-compressed files. Keep
GZIP as the portable
+ // default across supported Java deployments.
+ return GZIP_COMPRESSION_CODEC;
+ }
+ }
+
+ private static Option<String> getSparkRuntimeVersion() {
+ try {
+ Class<?> sparkPackageClass = Class.forName("org.apache.spark.package$");
+ Object sparkPackage = sparkPackageClass.getField("MODULE$").get(null);
+ return Option.of((String)
sparkPackageClass.getMethod("SPARK_VERSION").invoke(sparkPackage));
+ } catch (ReflectiveOperationException | LinkageError e) {
+ log.debug("Unable to resolve the Spark runtime version; using the legacy
Parquet compression codec default: {}", e.toString());
+ return Option.empty();
+ }
+ }
+
/**
* base properties.
*/
@@ -3813,6 +3851,10 @@ public class HoodieWriteConfig extends HoodieConfig {
}
// Check for mandatory properties
writeConfig.setDefaults(HoodieWriteConfig.class.getName());
+ // Resolve this engine-dependent default only after the final engine
type is known. The
+ // property remains unset in a partial HoodieStorageConfig, so explicit
values are preserved.
+
writeConfig.setDefaultValue(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME,
+ getDefaultParquetCompressionCodec(engineType));
// Make sure the props is propagated
writeConfig.setDefaultOnCondition(
!isIndexConfigSet,
HoodieIndexConfig.newBuilder().withEngineType(engineType).fromProperties(
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieBinaryCopyHandle.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieBinaryCopyHandle.java
index e209acf232d9..e79faaec9291 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieBinaryCopyHandle.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieBinaryCopyHandle.java
@@ -19,7 +19,6 @@
package org.apache.hudi.io;
import org.apache.hudi.client.WriteStatus;
-import org.apache.hudi.common.config.HoodieStorageConfig;
import org.apache.hudi.common.engine.TaskContextSupplier;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.HoodieWriteStat;
@@ -104,7 +103,7 @@ public class HoodieBinaryCopyHandle<T, I, K, O> extends
HoodieWriteHandle<T, I,
writeStatus.setStat(new HoodieWriteStat());
this.writer = new HoodieParquetFileBinaryCopier(
conf,
-
CompressionCodecName.fromConf(config.getStringOrDefault(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME)),
+ CompressionCodecName.fromConf(config.getParquetCompressionCodec()),
fileMetadataMerger);
// From the TABLE config, as the write handles do: the mode is a table
property, and the copier
// otherwise rewrites _hoodie_file_name on a table that does not populate
it.
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieMetadataWriteUtils.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieMetadataWriteUtils.java
index 8ab4c12a7902..9ff6f1a06cc3 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieMetadataWriteUtils.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieMetadataWriteUtils.java
@@ -313,6 +313,7 @@ public class HoodieMetadataWriteUtils {
.withStorageConfig(HoodieStorageConfig.newBuilder().hfileMaxFileSize(MDT_MAX_HFILE_SIZE_BYTES)
.allowDuplicatesWithHfileWrites(writeConfig.allowDuplicatesWithHfileWrites())
.logFileMaxSize(maxLogFileSizeBytes)
+ .parquetCompressionCodec(writeConfig.getParquetCompressionCodec())
// Keeping the log blocks as large as the log files themselves
reduces the number of HFile blocks to be checked for
// presence of keys
.logFileDataBlockMaxSize(maxLogFileSizeBytes)
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfig.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfig.java
index 80e0e84d4f52..22cea245639d 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfig.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfig.java
@@ -156,6 +156,42 @@ public class TestHoodieWriteConfig {
EngineType.JAVA, HoodieIndex.IndexType.SIMPLE));
}
+ @Test
+ public void testDefaultParquetCompressionCodecAccordingToEngine() {
+ assertEquals("zstd",
HoodieWriteConfig.getDefaultParquetCompressionCodec(EngineType.FLINK));
+ assertEquals("gzip",
HoodieWriteConfig.getDefaultParquetCompressionCodec(EngineType.JAVA));
+
+ HoodieWriteConfig flinkConfig = HoodieWriteConfig.newBuilder()
+ .withEngineType(EngineType.FLINK)
+ .withPath("/tmp")
+
.withStorageConfig(HoodieStorageConfig.newBuilder().parquetWriteLegacyFormat("false").build())
+ .build();
+ assertEquals("zstd", flinkConfig.getParquetCompressionCodec());
+
+ HoodieWriteConfig javaConfig = HoodieWriteConfig.newBuilder()
+ .withEngineType(EngineType.JAVA)
+ .withPath("/tmp")
+
.withStorageConfig(HoodieStorageConfig.newBuilder().parquetWriteLegacyFormat("false").build())
+ .build();
+ assertEquals("gzip", javaConfig.getParquetCompressionCodec());
+
+ javaConfig = HoodieWriteConfig.newBuilder()
+ .withEngineType(EngineType.JAVA)
+ .withPath("/tmp")
+
.withStorageConfig(HoodieStorageConfig.newBuilder().parquetCompressionCodec("zstd").build())
+ .build();
+ assertEquals("zstd", javaConfig.getParquetCompressionCodec());
+
+ Properties explicitCodec = new Properties();
+
explicitCodec.setProperty(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME.key(),
"gzip");
+ flinkConfig = HoodieWriteConfig.newBuilder()
+ .withEngineType(EngineType.FLINK)
+ .withPath("/tmp")
+
.withStorageConfig(HoodieStorageConfig.newBuilder().fromProperties(explicitCodec).build())
+ .build();
+ assertEquals("gzip", flinkConfig.getParquetCompressionCodec());
+ }
+
@Test
public void testDefaultBulkInsertSortModeForLsmLayout() {
HoodieWriteConfig defaultLayoutConfig = HoodieWriteConfig.newBuilder()
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieMetadataWriteUtils.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieMetadataWriteUtils.java
index a68bd289119a..abc5e5a42a77 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieMetadataWriteUtils.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieMetadataWriteUtils.java
@@ -129,6 +129,20 @@ public class TestHoodieMetadataWriteUtils {
assertEquals(hfileBloomFilterEnabled,
metadataWriteConfig.hfileBloomFilterEnabled());
}
+ @Test
+ public void testCreateMetadataWriteConfigPropagatesParquetCompressionCodec()
{
+ HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+ .withPath("/tmp/base_path/")
+ .withStorageConfig(HoodieStorageConfig.newBuilder()
+ .parquetCompressionCodec("snappy")
+ .build())
+ .build();
+
+ HoodieWriteConfig metadataWriteConfig =
HoodieMetadataWriteUtils.createMetadataWriteConfig(
+ writeConfig, HoodieFailedWritesCleaningPolicy.EAGER,
HoodieTableVersion.EIGHT);
+ assertEquals("snappy", metadataWriteConfig.getParquetCompressionCodec());
+ }
+
@ParameterizedTest
@EnumSource(value = MetricsReporterType.class, names = {
"GRAPHITE", "JMX", "PROMETHEUS_PUSHGATEWAY", "M3", "PROMETHEUS"
diff --git
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataFileWriterFactory.java
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataFileWriterFactory.java
index b75aeaede82e..f6de482e1688 100644
---
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataFileWriterFactory.java
+++
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataFileWriterFactory.java
@@ -144,7 +144,7 @@ public class HoodieRowDataFileWriterFactory extends
HoodieFileWriterFactory {
HoodieConfig config, HoodieRowDataParquetWriteSupport writeSupport) {
return new HoodieParquetConfig<>(
writeSupport,
-
getCompressionCodecName(config.getStringOrDefault(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME)),
+
getCompressionCodecName(config.getString(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME)),
config.getIntOrDefault(HoodieStorageConfig.PARQUET_BLOCK_SIZE),
config.getIntOrDefault(HoodieStorageConfig.PARQUET_PAGE_SIZE),
config.getLongOrDefault(HoodieStorageConfig.PARQUET_MAX_FILE_SIZE),
diff --git
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowDataParquetConfigInjector.java
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowDataParquetConfigInjector.java
index d2e109fbb52f..9a3cffb11111 100644
---
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowDataParquetConfigInjector.java
+++
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowDataParquetConfigInjector.java
@@ -126,6 +126,7 @@ public class TestHoodieRowDataParquetConfigInjector extends
HoodieFlinkClientTes
// Create config with the custom injector
HoodieConfig config = new HoodieConfig();
+ config.setValue(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME,
"zstd");
config.setValue(HoodieStorageConfig.PARQUET_DICTIONARY_ENABLED, "true");
// Start with dictionary enabled
config.setValue(HoodieStorageConfig.HOODIE_PARQUET_CONFIG_INJECTOR_CLASS,
DisableDictionaryInjector.class.getName());
@@ -206,6 +207,7 @@ public class TestHoodieRowDataParquetConfigInjector extends
HoodieFlinkClientTes
// Create config WITHOUT injector - should use default settings
HoodieConfig config = new HoodieConfig();
+ config.setValue(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME,
"zstd");
config.setValue(HoodieStorageConfig.PARQUET_DICTIONARY_ENABLED, "true");
// Create writer and write some data
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkFileWriterFactory.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkFileWriterFactory.java
index 3e5e4d716e5a..bdd316adb586 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkFileWriterFactory.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkFileWriterFactory.java
@@ -105,7 +105,7 @@ public class HoodieSparkFileWriterFactory extends
HoodieFileWriterFactory {
StorageConfiguration storageConfiguration = injectedConfigs.getLeft();
HoodieConfig hoodieConfig = injectedConfigs.getRight();
- String compressionCodecName =
hoodieConfig.getStringOrDefault(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME);
+ String compressionCodecName =
hoodieConfig.getString(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME);
// Support PARQUET_COMPRESSION_CODEC_NAME is ""
if (compressionCodecName.isEmpty()) {
compressionCodecName = null;
@@ -130,7 +130,7 @@ public class HoodieSparkFileWriterFactory extends
HoodieFileWriterFactory {
HoodieSchema schema) throws
IOException {
boolean enableBloomFilter = false;
HoodieRowParquetWriteSupport writeSupport =
getHoodieRowParquetWriteSupport(storage.getConf(), schema, config,
enableBloomFilter);
- String compressionCodecName =
config.getStringOrDefault(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME);
+ String compressionCodecName =
config.getString(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME);
// Support PARQUET_COMPRESSION_CODEC_NAME is ""
if (compressionCodecName.isEmpty()) {
compressionCodecName = null;
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/config/HoodieStorageConfig.java
b/hudi-common/src/main/java/org/apache/hudi/common/config/HoodieStorageConfig.java
index 07263f156ca2..6f2fc87c6073 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/config/HoodieStorageConfig.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/config/HoodieStorageConfig.java
@@ -215,8 +215,12 @@ public class HoodieStorageConfig extends HoodieConfig {
// Default compression codec for parquet
public static final ConfigProperty<String> PARQUET_COMPRESSION_CODEC_NAME =
ConfigProperty
.key("hoodie.parquet.compression.codec")
- .defaultValue("gzip")
- .withDocumentation("Compression Codec for parquet files");
+ .noDefaultValue("ZSTD for Flink and Spark 3.5 or newer; GZIP for Java
and older Spark versions")
+ .withDocumentation("Compression codec for Parquet base and native log
files. The default is ZSTD for Flink "
+ + "and Spark 3.5 or newer, and GZIP for Java and older Spark
versions. Spark 3.3 and 3.4 use a "
+ + "non-vectorized file-group reader affected by PARQUET-2160 when
reading ZSTD files, which can leak "
+ + "off-heap memory; upgrade to Spark 3.5 or newer before using ZSTD.
An explicitly configured value "
+ + "always takes precedence over the engine default.");
public static final ConfigProperty<Boolean> PARQUET_DICTIONARY_ENABLED =
ConfigProperty
.key("hoodie.parquet.dictionary.enabled")
@@ -532,7 +536,7 @@ public class HoodieStorageConfig extends HoodieConfig {
* @deprecated Use {@link #PARQUET_COMPRESSION_CODEC_NAME} and its methods
instead
*/
@Deprecated
- public static final String DEFAULT_PARQUET_COMPRESSION_CODEC =
PARQUET_COMPRESSION_CODEC_NAME.defaultValue();
+ public static final String DEFAULT_PARQUET_COMPRESSION_CODEC = "zstd";
/**
* @deprecated Use {@link #HFILE_COMPRESSION_ALGORITHM_NAME} and its methods
instead
*/
diff --git
a/hudi-common/src/test/java/org/apache/hudi/common/config/TestHoodieStorageConfig.java
b/hudi-common/src/test/java/org/apache/hudi/common/config/TestHoodieStorageConfig.java
index a7901c3e3e0f..0ca167124ddf 100644
---
a/hudi-common/src/test/java/org/apache/hudi/common/config/TestHoodieStorageConfig.java
+++
b/hudi-common/src/test/java/org/apache/hudi/common/config/TestHoodieStorageConfig.java
@@ -30,6 +30,7 @@ import static
org.apache.hudi.common.config.HoodieStorageConfig.BLOOM_FILTER_FPP
import static
org.apache.hudi.common.config.HoodieStorageConfig.BLOOM_FILTER_NUM_ENTRIES_VALUE;
import static
org.apache.hudi.common.config.HoodieStorageConfig.BLOOM_FILTER_TYPE;
import static
org.apache.hudi.common.config.HoodieStorageConfig.HFILE_WITH_BLOOM_FILTER_ENABLED;
+import static
org.apache.hudi.common.config.HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
@@ -78,6 +79,21 @@ public class TestHoodieStorageConfig {
assertEquals(BLOOM_FILTER_DYNAMIC_MAX_ENTRIES.defaultValue(),
storageConfig.getString(BLOOM_FILTER_DYNAMIC_MAX_ENTRIES));
}
+ @Test
+ void testParquetCompressionCodecDefaultIsDeferred() {
+ HoodieStorageConfig defaultStorageConfig =
HoodieStorageConfig.newBuilder().build();
+ assertFalse(PARQUET_COMPRESSION_CODEC_NAME.hasDefaultValue());
+ assertEquals("ZSTD for Flink and Spark 3.5 or newer; GZIP for Java and
older Spark versions",
+ PARQUET_COMPRESSION_CODEC_NAME.getDocOnDefaultValue());
+ assertFalse(defaultStorageConfig.contains(PARQUET_COMPRESSION_CODEC_NAME));
+
+ HoodieStorageConfig explicitStorageConfig =
HoodieStorageConfig.newBuilder()
+ .parquetCompressionCodec("zstd")
+ .build();
+ assertTrue(explicitStorageConfig.contains(PARQUET_COMPRESSION_CODEC_NAME));
+ assertEquals("zstd",
explicitStorageConfig.getString(PARQUET_COMPRESSION_CODEC_NAME));
+ }
+
@Test
void testHFileBloomFilterBuilder() {
HoodieStorageConfig defaultStorageConfig =
HoodieStorageConfig.newBuilder().build();
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/utils/TestFlinkWriteClients.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/utils/TestFlinkWriteClients.java
index f0d0c7968bf4..fba0e2ee098b 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/utils/TestFlinkWriteClients.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/utils/TestFlinkWriteClients.java
@@ -27,6 +27,7 @@ import
org.apache.hudi.client.model.PartialUpdateFlinkRecordMerger;
import org.apache.hudi.client.transaction.lock.FileSystemBasedLockProvider;
import org.apache.hudi.common.bloom.BloomFilterTypeCode;
import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.config.HoodieStorageConfig;
import org.apache.hudi.common.config.RecordMergeMode;
import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
import org.apache.hudi.common.model.EventTimeAvroPayload;
@@ -80,6 +81,16 @@ public class TestFlinkWriteClients {
this.conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath());
}
+ @Test
+ void testParquetCompressionCodecDefaultAndOverride() {
+ HoodieWriteConfig writeConfig =
FlinkWriteClients.getHoodieClientConfig(conf, false, false);
+ assertEquals("zstd", writeConfig.getParquetCompressionCodec());
+
+ conf.setString(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME.key(),
"gzip");
+ writeConfig = FlinkWriteClients.getHoodieClientConfig(conf, false, false);
+ assertEquals("gzip", writeConfig.getParquetCompressionCodec());
+ }
+
@Test
void testAutoSetupLockProvider() throws Exception {
conf.set(FlinkOptions.METADATA_ENABLED, true);
diff --git
a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroFileWriterFactory.java
b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroFileWriterFactory.java
index 9c91068a38e5..08deb5841d7e 100644
---
a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroFileWriterFactory.java
+++
b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroFileWriterFactory.java
@@ -117,7 +117,7 @@ public class HoodieAvroFileWriterFactory extends
HoodieFileWriterFactory {
HoodieAvroWriteSupport writeSupport = getHoodieAvroWriteSupport(schema,
hoodieConfig, storageConfiguration, enableBloomFilter(metaFieldsMode,
hoodieConfig));
- String compressionCodecName =
hoodieConfig.getStringOrDefault(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME);
+ String compressionCodecName =
hoodieConfig.getString(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME);
// Support PARQUET_COMPRESSION_CODEC_NAME is ""
if (compressionCodecName.isEmpty()) {
compressionCodecName = null;
diff --git
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/testutils/HoodieCommonTestHarness.java
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/testutils/HoodieCommonTestHarness.java
index ed6a6be4105f..891487ea51d0 100644
---
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/testutils/HoodieCommonTestHarness.java
+++
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/testutils/HoodieCommonTestHarness.java
@@ -75,7 +75,6 @@ import java.util.function.Predicate;
import java.util.stream.Collectors;
import static
org.apache.hudi.common.config.HoodieStorageConfig.HFILE_COMPRESSION_ALGORITHM_NAME;
-import static
org.apache.hudi.common.config.HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME;
/**
* The common hoodie test harness to provide the basic infrastructure.
@@ -413,7 +412,7 @@ public class HoodieCommonTestHarness {
case HFILE_DATA_BLOCK:
return new HoodieHFileDataBlock(records, header,
HFILE_COMPRESSION_ALGORITHM_NAME.defaultValue(), pathForReader);
case PARQUET_DATA_BLOCK:
- return new HoodieParquetDataBlock(records, header,
HoodieRecord.RECORD_KEY_METADATA_FIELD,
PARQUET_COMPRESSION_CODEC_NAME.defaultValue(), 0.1, true);
+ return new HoodieParquetDataBlock(records, header,
HoodieRecord.RECORD_KEY_METADATA_FIELD, "zstd", 0.1, true);
default:
throw new RuntimeException("Unknown data block type " + dataBlockType);
}
diff --git
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/testutils/reader/HoodieFileSliceTestUtils.java
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/testutils/reader/HoodieFileSliceTestUtils.java
index aa0e606bb9e9..7d875aacb7be 100644
---
a/hudi-hadoop-common/src/test/java/org/apache/hudi/common/testutils/reader/HoodieFileSliceTestUtils.java
+++
b/hudi-hadoop-common/src/test/java/org/apache/hudi/common/testutils/reader/HoodieFileSliceTestUtils.java
@@ -81,7 +81,6 @@ import java.util.stream.Collectors;
import java.util.stream.IntStream;
import static
org.apache.hudi.common.config.HoodieStorageConfig.HFILE_COMPRESSION_ALGORITHM_NAME;
-import static
org.apache.hudi.common.config.HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME;
import static
org.apache.hudi.common.table.log.block.HoodieLogBlock.HeaderMetadataType.BASE_FILE_INSTANT_TIME_OF_RECORD_POSITIONS;
import static
org.apache.hudi.common.table.log.block.HoodieLogBlock.HoodieLogBlockType.DELETE_BLOCK;
import static
org.apache.hudi.common.table.log.block.HoodieLogBlock.HoodieLogBlockType.PARQUET_DATA_BLOCK;
@@ -229,7 +228,7 @@ public class HoodieFileSliceTestUtils {
records,
header,
HoodieRecord.RECORD_KEY_METADATA_FIELD,
- PARQUET_COMPRESSION_CODEC_NAME.defaultValue(),
+ "zstd",
0.1,
true);
default:
diff --git
a/hudi-hadoop-common/src/test/java/org/apache/hudi/io/storage/hadoop/TestHoodieAvroFileWriterFactoryVariantInference.java
b/hudi-hadoop-common/src/test/java/org/apache/hudi/io/storage/hadoop/TestHoodieAvroFileWriterFactoryVariantInference.java
index 380868c118df..161f537dd101 100644
---
a/hudi-hadoop-common/src/test/java/org/apache/hudi/io/storage/hadoop/TestHoodieAvroFileWriterFactoryVariantInference.java
+++
b/hudi-hadoop-common/src/test/java/org/apache/hudi/io/storage/hadoop/TestHoodieAvroFileWriterFactoryVariantInference.java
@@ -73,6 +73,7 @@ public class TestHoodieAvroFileWriterFactoryVariantInference {
HoodieSchemaField.of("v",
HoodieSchema.createNullable(HoodieSchema.createVariant()))));
HoodieConfig config = new HoodieConfig();
config.setValue(HoodieStorageConfig.PARQUET_VARIANT_SHREDDING_SCHEMA_INFERENCE_ENABLED,
"true");
+ config.setValue(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME,
"zstd");
// Name a provider explicitly: the factory also declines when no shredding
provider is available,
// and this module ships none, so without this the inferrer gate (the one
under test) would never
// be reached. The class is never loaded: the write support only resolves
it for shredded schemas.
diff --git
a/hudi-hadoop-common/src/test/java/org/apache/hudi/io/storage/hadoop/TestHoodieAvroParquetConfigInjector.java
b/hudi-hadoop-common/src/test/java/org/apache/hudi/io/storage/hadoop/TestHoodieAvroParquetConfigInjector.java
index 4748b3db5032..48203238b9b9 100644
---
a/hudi-hadoop-common/src/test/java/org/apache/hudi/io/storage/hadoop/TestHoodieAvroParquetConfigInjector.java
+++
b/hudi-hadoop-common/src/test/java/org/apache/hudi/io/storage/hadoop/TestHoodieAvroParquetConfigInjector.java
@@ -36,6 +36,7 @@ import org.apache.parquet.column.Encoding;
import org.apache.parquet.hadoop.ParquetFileReader;
import org.apache.parquet.hadoop.metadata.BlockMetaData;
import org.apache.parquet.hadoop.metadata.ColumnChunkMetaData;
+import org.apache.parquet.hadoop.metadata.CompressionCodecName;
import org.apache.parquet.hadoop.metadata.ParquetMetadata;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -43,6 +44,7 @@ import org.junit.jupiter.api.io.TempDir;
import java.io.IOException;
import java.util.List;
+import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
@@ -70,6 +72,7 @@ public class TestHoodieAvroParquetConfigInjector {
// Create config with the custom injector
HoodieConfig config = new HoodieConfig();
+ config.setValue(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME,
"zstd");
config.setValue(HoodieStorageConfig.PARQUET_DICTIONARY_ENABLED, "true");
// Start with dictionary enabled
config.setValue(HoodieStorageConfig.HOODIE_PARQUET_CONFIG_INJECTOR_CLASS,
DisableDictionaryInjector.class.getName());
@@ -148,6 +151,7 @@ public class TestHoodieAvroParquetConfigInjector {
// Create config WITHOUT injector - should use default settings
HoodieConfig config = new HoodieConfig();
+ config.setValue(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME,
"zstd");
config.setValue(HoodieStorageConfig.PARQUET_DICTIONARY_ENABLED, "true");
// Create writer and write some data
@@ -165,5 +169,10 @@ public class TestHoodieAvroParquetConfigInjector {
// Verify the parquet file was created
assertTrue(storage.exists(parquetPath));
+
+ try (ParquetFileReader reader = ParquetFileReader.open(new
Configuration(), new Path(parquetPath.toUri()))) {
+ assertEquals(CompressionCodecName.ZSTD,
+
reader.getFooter().getBlocks().get(0).getColumns().get(0).getCodec());
+ }
}
}
diff --git
a/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
b/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
index 24eced21dca6..769af07c7388 100644
---
a/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
+++
b/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
@@ -22,7 +22,9 @@ package org.apache.hudi.metadata;
import org.apache.hudi.avro.model.HoodieInstantInfo;
import org.apache.hudi.avro.model.HoodieRollbackMetadata;
import org.apache.hudi.avro.model.HoodieRollbackPlan;
+import org.apache.hudi.common.config.HoodieConfig;
import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.config.HoodieStorageConfig;
import org.apache.hudi.common.data.HoodieData;
import org.apache.hudi.common.engine.HoodieLocalEngineContext;
import org.apache.hudi.common.function.SerializableBiFunction;
@@ -335,11 +337,16 @@ public class TestHoodieTableMetadataUtil extends
HoodieCommonTestHarness {
List<HoodieRecord> records,
HoodieTableMetaClient metaClient,
HoodieLocalEngineContext engineContext)
throws IOException {
+ HoodieConfig writerConfig = new HoodieConfig();
+ writerConfig.getProps().putAll(metaClient.getTableConfig().getProps());
+ // This low-level test path bypasses HoodieWriteConfig, and does not
contain codec configuration,
+ // so set the codec explicitly here.
+ writerConfig.setValue(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME,
"zstd");
HoodieFileWriter writer = HoodieFileWriterFactory.getFileWriter(
instant,
path,
metaClient.getStorage(),
- metaClient.getTableConfig(),
+ writerConfig,
HOODIE_SCHEMA_WITH_METADATA_FIELDS,
engineContext.getTaskContextSupplier(),
HoodieRecord.HoodieRecordType.AVRO);
diff --git
a/hudi-hadoop-mr/src/test/java/org/apache/hudi/hadoop/testutils/InputFormatTestUtil.java
b/hudi-hadoop-mr/src/test/java/org/apache/hudi/hadoop/testutils/InputFormatTestUtil.java
index c7d363a0cef6..0ea8ed524761 100644
---
a/hudi-hadoop-mr/src/test/java/org/apache/hudi/hadoop/testutils/InputFormatTestUtil.java
+++
b/hudi-hadoop-mr/src/test/java/org/apache/hudi/hadoop/testutils/InputFormatTestUtil.java
@@ -85,7 +85,6 @@ import java.util.UUID;
import java.util.stream.Collectors;
import static
org.apache.hudi.common.config.HoodieStorageConfig.HFILE_COMPRESSION_ALGORITHM_NAME;
-import static
org.apache.hudi.common.config.HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME;
public class InputFormatTestUtil {
@@ -445,7 +444,7 @@ public class InputFormatTestUtil {
hoodieRecords, header,
HFILE_COMPRESSION_ALGORITHM_NAME.defaultValue(), writer.getLogFile().getPath());
} else if (logBlockType ==
HoodieLogBlock.HoodieLogBlockType.PARQUET_DATA_BLOCK) {
dataBlock = new HoodieParquetDataBlock(hoodieRecords, header,
- HoodieRecord.RECORD_KEY_METADATA_FIELD,
PARQUET_COMPRESSION_CODEC_NAME.defaultValue(), 0.1, true);
+ HoodieRecord.RECORD_KEY_METADATA_FIELD, "zstd", 0.1, true);
} else {
dataBlock = new HoodieAvroDataBlock(
hoodieRecords, header, HoodieRecord.RECORD_KEY_METADATA_FIELD);
diff --git a/hudi-spark-datasource/README.md b/hudi-spark-datasource/README.md
index 1a96eafa0a01..413eeaf33e0d 100644
--- a/hudi-spark-datasource/README.md
+++ b/hudi-spark-datasource/README.md
@@ -50,6 +50,16 @@ The modules are organized in a layered architecture to
maximize code reuse acros
| 4.1.x | `hudi-spark4.1.x` | 2.13 | 17+ | `-Dspark4.1` |
| 4.2.x | `hudi-spark4.2.x` | 2.13 | 17+ | `-Dspark4.2` |
+## Parquet Compression Compatibility
+
+Hudi defaults `hoodie.parquet.compression.codec` to ZSTD on Spark 3.5 and
newer and to GZIP on
+Spark 3.3 and 3.4. The non-vectorized Hudi file-group reader in Spark 3.3 and
3.4 uses parquet-java
+1.12.x, which is affected by
[PARQUET-2160](https://issues.apache.org/jira/browse/PARQUET-2160)
+and can leak off-heap memory while reading ZSTD files. Upgrade to Spark 3.5 or
newer before using
+ZSTD. The GZIP write default on older Spark versions does not eliminate the
risk of reading ZSTD
+files produced by another engine. Explicitly setting
`hoodie.parquet.compression.codec` overrides
+the version-specific default.
+
## Key Features
- **DataSource V1 Support**: Full integration with Spark's DataSource API
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestDataSourceUtils.java
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestDataSourceUtils.java
index 840156cde89c..8755120bb641 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestDataSourceUtils.java
+++
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestDataSourceUtils.java
@@ -21,6 +21,7 @@ package org.apache.hudi;
import org.apache.hudi.client.SparkRDDWriteClient;
import org.apache.hudi.client.WriteStatus;
import org.apache.hudi.common.avro.HoodieAvroUtils;
+import org.apache.hudi.common.config.HoodieStorageConfig;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
import org.apache.hudi.common.model.HoodieRecord;
@@ -29,6 +30,7 @@ import org.apache.hudi.common.model.WriteOperationType;
import org.apache.hudi.common.schema.HoodieSchema;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.SerializationUtils;
+import org.apache.hudi.common.util.StringUtils;
import org.apache.hudi.common.util.collection.ImmutablePair;
import org.apache.hudi.config.HoodieClusteringConfig;
import org.apache.hudi.config.HoodieWriteConfig;
@@ -120,6 +122,30 @@ public class TestDataSourceUtils extends
HoodieClientTestBase {
config = HoodieWriteConfig.newBuilder().withPath("/").build();
}
+ @Test
+ public void testSparkVersionSpecificParquetCompressionCodecDefault() {
+ String expectedCodec =
StringUtils.compareVersions(HoodieSparkUtils.getSparkVersion(), "3.5.0") >= 0
+ ? "zstd" : "gzip";
+ assertEquals(expectedCodec, config.getParquetCompressionCodec());
+
+ HoodieWriteConfig configWithPartialStorage = HoodieWriteConfig.newBuilder()
+ .withPath("/")
+
.withStorageConfig(HoodieStorageConfig.newBuilder().parquetWriteLegacyFormat("false").build())
+ .build();
+ assertEquals(expectedCodec,
configWithPartialStorage.getParquetCompressionCodec());
+
+ Map<String, String> params = new HashMap<>();
+ params.put(DataSourceWriteOptions.TABLE_TYPE().key(),
DataSourceWriteOptions.COW_TABLE_TYPE_OPT_VAL());
+ HoodieWriteConfig dataSourceConfig = DataSourceUtils.createHoodieConfig(
+ avroSchemaString, config.getBasePath(), "test", params);
+ assertEquals(expectedCodec, dataSourceConfig.getParquetCompressionCodec());
+
+ params.put(HoodieStorageConfig.PARQUET_COMPRESSION_CODEC_NAME.key(),
"snappy");
+ dataSourceConfig = DataSourceUtils.createHoodieConfig(
+ avroSchemaString, config.getBasePath(), "test", params);
+ assertEquals("snappy", dataSourceConfig.getParquetCompressionCodec());
+ }
+
@Test
public void testAvroRecordsFieldConversion() {
diff --git a/packaging/bundle-validation/spark_hadoop_mr/write.scala
b/packaging/bundle-validation/spark_hadoop_mr/write.scala
index e36ccc120373..e47ca4fbc7b6 100644
--- a/packaging/bundle-validation/spark_hadoop_mr/write.scala
+++ b/packaging/bundle-validation/spark_hadoop_mr/write.scala
@@ -31,8 +31,12 @@ val basePath = "file:///tmp/hudi-bundles/tests/" + tableName
val dataGen = new DataGenerator
val inserts =
convertToStringList(dataGen.generateInserts(expected)).asScala.toSeq
val df = spark.read.json(spark.sparkContext.parallelize(inserts, 2))
+// Hive 3.1.3 uses Hadoop's native ZSTD codec to read Parquet files, but the
Alpine-based
+// bundle-validation environment cannot load the native Hadoop library. Keep
this Spark-to-Hive
+// compatibility test on GZIP; engine-specific ZSTD defaults are validated
separately.
df.write.format("hudi").
options(getQuickstartWriteConfigs).
+ option("hoodie.parquet.compression.codec", "gzip").
option(PRECOMBINE_FIELD_OPT_KEY, "ts").
option(RECORDKEY_FIELD_OPT_KEY, "uuid").
option(PARTITIONPATH_FIELD_OPT_KEY, "partitionpath").
diff --git a/packaging/hudi-flink-bundle/pom.xml
b/packaging/hudi-flink-bundle/pom.xml
index 25e7529ef2f7..2c6eb05bc8ae 100644
--- a/packaging/hudi-flink-bundle/pom.xml
+++ b/packaging/hudi-flink-bundle/pom.xml
@@ -69,7 +69,8 @@
</transformer>
<transformer
implementation="org.apache.maven.plugins.shade.resource.IncludeResourceTransformer">
<resource>META-INF/LICENSE</resource>
- <file>target/classes/META-INF/LICENSE</file>
+ <!-- zstd-jni does not include its BSD license as a JAR
resource. -->
+ <file>${project.basedir}/../../LICENSE</file>
</transformer>
<transformer
implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/>
</transformers>
@@ -105,6 +106,7 @@
<include>org.apache.parquet:parquet-format-structures</include>
<include>org.apache.parquet:parquet-encoding</include>
<include>org.apache.parquet:parquet-jackson</include>
+ <include>com.github.luben:zstd-jni</include>
<include>org.apache.avro:avro</include>
<include>joda-time:joda-time</include>