danny0405 commented on code in PR #19685:
URL: https://github.com/apache/hudi/pull/19685#discussion_r3870354454


##########
hudi-common/src/main/java/org/apache/hudi/common/config/HoodieStorageConfig.java:
##########
@@ -534,12 +549,47 @@ public String getBloomFilterType() {
     return getString(BLOOM_FILTER_TYPE);
   }
 
+  @VisibleForTesting
+  static String getDefaultParquetCompressionCodec(EngineType engineType) {
+    switch (engineType) {
+      case FLINK:
+        return ZSTD_COMPRESSION_CODEC;
+      case SPARK:
+        // Uses ZSTD as the default Parquet compression codec with Spark 3.5 
and newer. Spark 3.3 and 3.4
+        // retain GZIP because the 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();
+    }
+  }
+
   public static HoodieStorageConfig.Builder newBuilder() {
     return new Builder();
   }
 
   public static class Builder {
 
+    private EngineType engineType = EngineType.SPARK;
     private final HoodieStorageConfig storageConfig = new 
HoodieStorageConfig();

Review Comment:
   The nested storage builder eagerly materializes a Spark default here, so a 
Java write config becomes classpath-sensitive. For example, in a JVM that also 
contains Spark 3.5 jars, 
HoodieWriteConfig.newBuilder().withEngineType(EngineType.JAVA).withStorageConfig(HoodieStorageConfig.newBuilder().parquetWriteLegacyFormat("false").build()).build()
 stores zstd in the inner config; withStorageConfig then marks storage as set, 
and the outer Java engine cannot restore the required gzip default. This 
affects existing callers that provide any partial storage config. Please defer 
this default until the outer write builder knows the engine, or otherwise 
preserve whether the codec was explicitly set.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to