This is an automated email from the ASF dual-hosted git repository.

voonhous 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 f6590309b42a fix: fix the file size calibration for native log format 
(#19293)
f6590309b42a is described below

commit f6590309b42ae4316e21d3a0b49d6bb1eef4e29b
Author: Danny Chan <[email protected]>
AuthorDate: Fri Jul 17 16:09:13 2026 +0800

    fix: fix the file size calibration for native log format (#19293)
    
    * fix: fix the file size calibration for native log format
    
    * fix variable name
    
    * remove the check for parquet log block since we have native log now
    
    * fix log file size calibration by block format
    
    * address Spark file sizing review comments
---
 .../SparkUpsertDeltaCommitPartitioner.java         |  6 ++--
 .../org/apache/hudi/common/model/FileSlice.java    | 39 +++++++++++++++++-----
 .../apache/hudi/common/model/TestFileSlice.java    | 24 +++++++++++++
 .../partitioner/profile/DeltaWriteProfile.java     | 15 ++++++---
 .../hudi/sink/partitioner/TestBucketAssigner.java  | 34 +++++++++++++++++++
 .../scala/org/apache/hudi/HoodieFileIndex.scala    |  7 ++--
 .../apache/hudi/PartitionDirectoryConverter.scala  |  8 ++---
 .../hudi/TestPartitionDirectoryConverter.scala     | 34 +++++++++++++++++--
 8 files changed, 142 insertions(+), 25 deletions(-)

diff --git 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/deltacommit/SparkUpsertDeltaCommitPartitioner.java
 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/deltacommit/SparkUpsertDeltaCommitPartitioner.java
index ad62f4e5a0e6..8319126cf3fd 100644
--- 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/deltacommit/SparkUpsertDeltaCommitPartitioner.java
+++ 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/action/deltacommit/SparkUpsertDeltaCommitPartitioner.java
@@ -73,13 +73,13 @@ public class SparkUpsertDeltaCommitPartitioner<T> extends 
UpsertPartitioner<T> {
       if (smallFileSlice.getBaseFile().isPresent()) {
         HoodieBaseFile baseFile = smallFileSlice.getBaseFile().get();
         sf.location = new HoodieRecordLocation(baseFile.getCommitTime(), 
baseFile.getFileId());
-        sf.sizeBytes = 
smallFileSlice.getTotalFileSizeAsParquetFormat(config.getLogFileToParquetCompressionRatio());
+        sf.sizeBytes = smallFileSlice.getTotalFileSizeAsParquetFormat(config);
         smallFileLocations.add(sf);
       } else {
         HoodieLogFile logFile = smallFileSlice.getLogFiles().findFirst().get();
         sf.location = new HoodieRecordLocation(logFile.getDeltaCommitTime(),
             logFile.getFileId());
-        sf.sizeBytes = 
smallFileSlice.getTotalFileSizeAsParquetFormat(config.getLogFileToParquetCompressionRatio());
+        sf.sizeBytes = smallFileSlice.getTotalFileSizeAsParquetFormat(config);
         smallFileLocations.add(sf);
       }
     }
@@ -122,7 +122,7 @@ public class SparkUpsertDeltaCommitPartitioner<T> extends 
UpsertPartitioner<T> {
   }
 
   private boolean isSmallFile(FileSlice fileSlice) {
-    long totalSize = 
fileSlice.getTotalFileSizeAsParquetFormat(config.getLogFileToParquetCompressionRatio());
+    long totalSize = fileSlice.getTotalFileSizeAsParquetFormat(config);
     return totalSize < config.getParquetMaxFileSize();
   }
 
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/model/FileSlice.java 
b/hudi-common/src/main/java/org/apache/hudi/common/model/FileSlice.java
index cf648b11fbfa..b1c625822536 100644
--- a/hudi-common/src/main/java/org/apache/hudi/common/model/FileSlice.java
+++ b/hudi-common/src/main/java/org/apache/hudi/common/model/FileSlice.java
@@ -18,7 +18,10 @@
 
 package org.apache.hudi.common.model;
 
+import org.apache.hudi.common.config.HoodieConfig;
+import org.apache.hudi.common.config.HoodieStorageConfig;
 import org.apache.hudi.common.function.SerializableFunctionUnchecked;
+import 
org.apache.hudi.common.table.log.block.HoodieLogBlock.HoodieLogBlockType;
 import org.apache.hudi.common.table.timeline.InstantComparison;
 import org.apache.hudi.common.util.Option;
 
@@ -193,24 +196,44 @@ public class FileSlice implements Serializable {
   }
 
   /**
-   * Get the total file size of a file slice similar on the base file.
-   * For the log file, we need to convert its size to the estimated size 
similar on the base file in a certain proportion
+   * Gets the estimated total size of a file slice in the base Parquet format.
+   *
+   * <p>Only inline Avro log files need size calibration. Native log files and 
inline non-Avro log files
+   * use their physical sizes directly.</p>
+   *
+   * @param writeConfig write config containing the log data block format and 
log-to-Parquet compression ratio
    */
-  public long getTotalFileSizeAsParquetFormat(double logFileFraction) {
-    long logFileSize = 
convertLogFilesSizeToExpectedParquetSize(logFileFraction);
+  public long getTotalFileSizeAsParquetFormat(HoodieConfig writeConfig) {
+    long logFileSize = convertLogFilesSizeToExpectedParquetSize(writeConfig);
     return getBaseFile().isPresent() ? getBaseFile().get().getFileSize() + 
logFileSize : logFileSize;
   }
 
-  private long convertLogFilesSizeToExpectedParquetSize(double 
logFileFraction) {
-    long totalSizeOfLogFiles =
+  private long convertLogFilesSizeToExpectedParquetSize(HoodieConfig 
writeConfig) {
+    long totalSizeOfInlineLogFiles =
+        logFiles.stream()
+            .filter(logFile -> !logFile.isNativeLogFile())
+            .map(HoodieLogFile::getFileSize)
+            .filter(size -> size > 0)
+            .reduce(Long::sum)
+            .orElse(0L);
+    long totalSizeOfNativeLogFiles =
         logFiles.stream()
+            .filter(HoodieLogFile::isNativeLogFile)
             .map(HoodieLogFile::getFileSize)
             .filter(size -> size > 0)
             .reduce(Long::sum)
             .orElse(0L);
-    // Here we assume that if there is no base parquet file, all log files 
contain only inserts.
+    // Here we assume that if there is no base parquet file, all calibrated 
log files contain only inserts.
     // We can then just get the parquet equivalent size of these log files, 
compare that with
     // {@link config.getParquetMaxFileSize()} and decide if there is scope to 
insert more rows
-    return (long) (totalSizeOfLogFiles * logFileFraction);
+    // The effective inline block type cannot be inferred from the log file 
name, and hudi-common does not have access
+    // to the table config fallback used by write clients. Legacy callers with 
no explicit block format default to Avro.
+    boolean isAvroDataBlocks = 
HoodieLogBlockType.fromId(writeConfig.getStringOrDefault(
+        HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT, "avro")) == 
HoodieLogBlockType.AVRO_DATA_BLOCK;
+    long calibratedInlineLogFileSize = isAvroDataBlocks
+        ? (long) (totalSizeOfInlineLogFiles
+            * 
writeConfig.getDoubleOrDefault(HoodieStorageConfig.LOGFILE_TO_PARQUET_COMPRESSION_RATIO_FRACTION))
+        : totalSizeOfInlineLogFiles;
+    return calibratedInlineLogFileSize + totalSizeOfNativeLogFiles;
   }
 }
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/model/TestFileSlice.java 
b/hudi-common/src/test/java/org/apache/hudi/common/model/TestFileSlice.java
index daf02bd0edd1..529d6742a1c3 100644
--- a/hudi-common/src/test/java/org/apache/hudi/common/model/TestFileSlice.java
+++ b/hudi-common/src/test/java/org/apache/hudi/common/model/TestFileSlice.java
@@ -18,6 +18,8 @@
 
 package org.apache.hudi.common.model;
 
+import org.apache.hudi.common.config.HoodieConfig;
+import org.apache.hudi.common.config.HoodieStorageConfig;
 import org.apache.hudi.storage.StoragePath;
 
 import org.junit.jupiter.api.Test;
@@ -164,4 +166,26 @@ public class TestFileSlice {
     assertEquals(fileSlice.getBaseFile(), 
fileSliceAfterFilterPart.getBaseFile());
     assertEquals(fileSlice.getBaseInstantTime(), 
fileSliceAfterFilterPart.getBaseInstantTime());
   }
+
+  @Test
+  void testGetTotalFileSizeAsParquetFormat() {
+    FileSlice fileSlice = new FileSlice(PARTITION_PATH, BASE_INSTANT, 
"file-id");
+    fileSlice.addLogFile(new HoodieLogFile(
+        new StoragePath(PARTITION_PATH + "/.file-id_002.log.1_1-0-1"), 1000L));
+    fileSlice.addLogFile(new HoodieLogFile(
+        new StoragePath(PARTITION_PATH + "/file-id_1-0-1_002_2.log.parquet"), 
2000L));
+
+    assertEquals(2100L, 
fileSlice.getTotalFileSizeAsParquetFormat(createLogConfig("avro", 0.1)));
+    assertEquals(3000L, 
fileSlice.getTotalFileSizeAsParquetFormat(createLogConfig("parquet", 0.1)));
+    assertEquals(3000L, 
fileSlice.getTotalFileSizeAsParquetFormat(createLogConfig("hfile", 0.1)));
+  }
+
+  private static HoodieConfig createLogConfig(String logDataBlockFormat, 
double logFileFraction) {
+    HoodieConfig config = new HoodieConfig();
+    config.setValue(HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT, 
logDataBlockFormat);
+    config.setValue(
+        HoodieStorageConfig.LOGFILE_TO_PARQUET_COMPRESSION_RATIO_FRACTION,
+        String.valueOf(logFileFraction));
+    return config;
+  }
 }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/profile/DeltaWriteProfile.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/profile/DeltaWriteProfile.java
index 2cdee3453a46..8f1e5c1d5e9f 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/profile/DeltaWriteProfile.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/partitioner/profile/DeltaWriteProfile.java
@@ -22,12 +22,14 @@ import 
org.apache.hudi.client.common.HoodieFlinkEngineContext;
 import org.apache.hudi.common.model.FileSlice;
 import org.apache.hudi.common.model.HoodieBaseFile;
 import org.apache.hudi.common.model.HoodieRecordLocation;
+import org.apache.hudi.common.table.HoodieTableVersion;
 import org.apache.hudi.common.table.log.block.HoodieLogBlock;
 import org.apache.hudi.common.table.timeline.HoodieInstant;
 import org.apache.hudi.common.table.timeline.HoodieTimeline;
 import org.apache.hudi.common.table.view.SyncableFileSystemView;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.table.action.commit.SmallFile;
+import org.apache.hudi.util.CommonClientUtils;
 
 import java.util.ArrayList;
 import java.util.List;
@@ -114,7 +116,7 @@ public class DeltaWriteProfile extends WriteProfile {
   }
 
   private long getTotalFileSize(FileSlice fileSlice) {
-    return 
fileSlice.getTotalFileSizeAsParquetFormat(logFileToParquetCompressionRatio());
+    return fileSlice.getTotalFileSizeAsParquetFormat(config);
   }
 
   private boolean isSmallFile(FileSlice fileSlice) {
@@ -123,10 +125,13 @@ public class DeltaWriteProfile extends WriteProfile {
   }
 
   private double logFileToParquetCompressionRatio() {
-    if (config.getLogDataBlockFormat().isPresent()
-        && config.getLogDataBlockFormat().get() == 
HoodieLogBlock.HoodieLogBlockType.PARQUET_DATA_BLOCK) {
-      return 1D;
+    // Delta commit metadata does not identify native and inline log files 
separately. The write version is expected
+    // to be reconciled with the table version, so version 10 and above can 
use the native log size directly.
+    if (config.getWriteVersion().lesserThan(HoodieTableVersion.TEN)
+        && CommonClientUtils.getLogBlockType(config, 
metaClient.getTableConfig())
+        == HoodieLogBlock.HoodieLogBlockType.AVRO_DATA_BLOCK) {
+      return config.getLogFileToParquetCompressionRatio();
     }
-    return config.getLogFileToParquetCompressionRatio();
+    return 1D;
   }
 }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestBucketAssigner.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestBucketAssigner.java
index 3b5c265f7c91..ee7bccf6a11c 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestBucketAssigner.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/partitioner/TestBucketAssigner.java
@@ -23,6 +23,7 @@ import org.apache.hudi.common.config.HoodieStorageConfig;
 import org.apache.hudi.common.model.HoodieCommitMetadata;
 import org.apache.hudi.common.model.HoodieRecordLocation;
 import org.apache.hudi.common.model.HoodieTableType;
+import org.apache.hudi.common.table.HoodieTableVersion;
 import org.apache.hudi.common.table.timeline.HoodieTimeline;
 import org.apache.hudi.config.HoodieCompactionConfig;
 import org.apache.hudi.config.HoodieWriteConfig;
@@ -476,6 +477,7 @@ public class TestBucketAssigner {
     
morConf.setString(HoodieCompactionConfig.COPY_ON_WRITE_RECORD_SIZE_ESTIMATE.key(),
 "1024");
     morConf.setString(HoodieCompactionConfig.PARQUET_SMALL_FILE_LIMIT.key(), 
"1");
     
morConf.setString(HoodieStorageConfig.LOGFILE_TO_PARQUET_COMPRESSION_RATIO_FRACTION.key(),
 "0.5");
+    morConf.set(FlinkOptions.WRITE_TABLE_VERSION, 
HoodieTableVersion.NINE.versionCode());
     StreamerUtil.initTableIfNotExists(morConf);
     TestData.writeData(TestData.DATA_SET_INSERT, morConf);
 
@@ -507,6 +509,7 @@ public class TestBucketAssigner {
     morConf.setString(HoodieCompactionConfig.PARQUET_SMALL_FILE_LIMIT.key(), 
"1");
     
morConf.setString(HoodieStorageConfig.LOGFILE_TO_PARQUET_COMPRESSION_RATIO_FRACTION.key(),
 "0.5");
     morConf.setString(HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key(), 
"parquet");
+    morConf.set(FlinkOptions.WRITE_TABLE_VERSION, 
HoodieTableVersion.NINE.versionCode());
     StreamerUtil.initTableIfNotExists(morConf);
     TestData.writeData(TestData.DATA_SET_INSERT, morConf);
 
@@ -528,6 +531,37 @@ public class TestBucketAssigner {
         writeProfile.getRecordsPerBucket(), 
is(morWriteConfig.getParquetMaxFileSize() / expectedAvgSize));
   }
 
+  @Test
+  public void 
testDeltaWriteProfileRecordsPerBucketSkipsCompressionRatioForNativeLogs() 
throws Exception {
+    File morPath = new File(tempFile, "mor_native_logs");
+    Configuration morConf = 
TestConfigurations.getDefaultConf(morPath.getAbsolutePath());
+    morConf.set(FlinkOptions.TABLE_TYPE, HoodieTableType.MERGE_ON_READ.name());
+    morConf.set(FlinkOptions.WRITE_PARQUET_MAX_FILE_SIZE, 1);
+    
morConf.setString(HoodieCompactionConfig.COPY_ON_WRITE_RECORD_SIZE_ESTIMATE.key(),
 "1024");
+    morConf.setString(HoodieCompactionConfig.PARQUET_SMALL_FILE_LIMIT.key(), 
"1");
+    
morConf.setString(HoodieStorageConfig.LOGFILE_TO_PARQUET_COMPRESSION_RATIO_FRACTION.key(),
 "0.5");
+    morConf.set(FlinkOptions.WRITE_TABLE_VERSION, 
HoodieTableVersion.TEN.versionCode());
+    StreamerUtil.initTableIfNotExists(morConf);
+    TestData.writeData(TestData.DATA_SET_INSERT, morConf);
+
+    HoodieWriteConfig morWriteConfig = 
FlinkWriteClients.getHoodieClientConfig(morConf);
+    HoodieFlinkEngineContext morContext = new HoodieFlinkEngineContext(
+        
HadoopFSUtils.getStorageConf(HadoopConfigurations.getHadoopConf(morConf)),
+        new FlinkTaskContextSupplier(null));
+
+    DeltaWriteProfile writeProfile = new DeltaWriteProfile(morWriteConfig, 
morContext);
+    String latestInstant = getLastCompleteInstant(writeProfile);
+    HoodieCommitMetadata commitMetadata = 
writeProfile.getMetadataCache().get(latestInstant);
+    assertNotNull(commitMetadata);
+    long expectedAvgSize = (long) Math.ceil(
+        1.0 * commitMetadata.fetchTotalBytesWritten() / 
commitMetadata.fetchTotalRecordsWritten());
+
+    assertThat("Average record size from native logs should not be corrected 
again",
+        writeProfile.getAvgSize(), is(expectedAvgSize));
+    assertThat("Records per bucket should use the uncorrected native log 
average record size",
+        writeProfile.getRecordsPerBucket(), 
is(morWriteConfig.getParquetMaxFileSize() / expectedAvgSize));
+  }
+
   @Test
   public void 
testDeltaWriteProfileReusesPreviousAvgSizeWhenNoEligibleDeltaCommitOnReload() 
throws Exception {
     File morPath = new File(tempFile, "mor_reuse_previous_avg");
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieFileIndex.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieFileIndex.scala
index 020fc3fc8c7b..cd10aeb9c1fe 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieFileIndex.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieFileIndex.scala
@@ -20,7 +20,7 @@ package org.apache.hudi
 import org.apache.hudi.DataSourceWriteOptions.{PARTITIONPATH_FIELD, 
RECORDKEY_FIELD}
 import org.apache.hudi.HoodieFileIndex.{collectReferencedColumns, 
convertFilterForTimestampKeyGenerator, getConfigProperties, 
DataSkippingFailureMode}
 import org.apache.hudi.HoodieSparkConfUtils.getConfigValue
-import org.apache.hudi.common.config.{HoodieMetadataConfig, TypedProperties}
+import org.apache.hudi.common.config.{HoodieConfig, HoodieMetadataConfig, 
TypedProperties}
 import 
org.apache.hudi.common.config.TimestampKeyGeneratorConfig.{TIMESTAMP_INPUT_DATE_FORMAT,
 TIMESTAMP_OUTPUT_DATE_FORMAT}
 import org.apache.hudi.common.model.{FileSlice, HoodieBaseFile, HoodieLogFile}
 import org.apache.hudi.common.table.{HoodieTableConfig, HoodieTableMetaClient}
@@ -106,6 +106,9 @@ case class HoodieFileIndex(spark: SparkSession,
 
   @transient protected var hasPushedDownPartitionPredicates: Boolean = false
 
+  @transient private lazy val hoodieConfig =
+    new HoodieConfig(TypedProperties.fromMap(options.filter(_._2 != 
null).asJava))
+
   /** True when any partition column is a nested field path (e.g. 
"nested_record.level"). */
   private val hasNestedPartitionColumns: Boolean =
     getPartitionColumns.exists(_.contains("."))
@@ -219,7 +222,7 @@ case class HoodieFileIndex(spark: SparkSession,
           PartitionDirectoryConverter.convertFileSliceToPartitionDirectory(
             partitionValues,
             fileSlice,
-            options)
+            hoodieConfig)
         } else {
           val baseFileStatusOpt = 
getBaseFileInfo(Option.apply(fileSlice.getBaseFile.orElse(null)))
           val logPathInfoStream = 
fileSlice.getLogFiles.map[StoragePathInfo](JFunction.toJavaFunction[HoodieLogFile,
 StoragePathInfo](lf => lf.getPathInfo))
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/PartitionDirectoryConverter.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/PartitionDirectoryConverter.scala
index ffd7bd1cf98c..2565278ab3eb 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/PartitionDirectoryConverter.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/PartitionDirectoryConverter.scala
@@ -17,7 +17,7 @@
 
 package org.apache.hudi
 
-import org.apache.hudi.common.config.HoodieStorageConfig
+import org.apache.hudi.common.config.HoodieConfig
 import org.apache.hudi.common.model.FileSlice
 
 import org.apache.hadoop.fs.{FileStatus, Path}
@@ -28,13 +28,11 @@ object PartitionDirectoryConverter extends 
SparkAdapterSupport {
 
   def convertFileSliceToPartitionDirectory(partitionValues: InternalRow,
                                            fileSlice: FileSlice,
-                                           options: Map[String, String]): 
PartitionDirectory = {
-    val logFileEstimationFraction = 
options.getOrElse(HoodieStorageConfig.LOGFILE_TO_PARQUET_COMPRESSION_RATIO_FRACTION.key(),
-      
HoodieStorageConfig.LOGFILE_TO_PARQUET_COMPRESSION_RATIO_FRACTION.defaultValue()).toDouble
+                                           config: HoodieConfig): 
PartitionDirectory = {
     // 1. Generate a delegate file for file slice, which spark uses to 
optimize rdd partition parallelism based on data such as file size
     //    - For file slice only has base file, we directly use the base file 
size as delegate file size
     //    - For file slice has log file, we estimate the delegate file size 
based on the log file size and option(base file) size
-    val estimationFileSize = 
fileSlice.getTotalFileSizeAsParquetFormat(logFileEstimationFraction)
+    val estimationFileSize = fileSlice.getTotalFileSizeAsParquetFormat(config)
     val fileInfo = if (fileSlice.getBaseFile.isPresent) {
       fileSlice.getBaseFile.get().getPathInfo
     } else {
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestPartitionDirectoryConverter.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestPartitionDirectoryConverter.scala
index b14921c2c098..204b2558960f 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestPartitionDirectoryConverter.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/TestPartitionDirectoryConverter.scala
@@ -17,7 +17,7 @@
 
 package org.apache.hudi
 
-import org.apache.hudi.common.config.HoodieStorageConfig
+import org.apache.hudi.common.config.{HoodieConfig, HoodieStorageConfig, 
TypedProperties}
 import org.apache.hudi.common.fs.FSUtils
 import org.apache.hudi.common.model.{FileSlice, HoodieBaseFile, 
HoodieFileGroupId, HoodieLogFile}
 import org.apache.hudi.storage.{StoragePath, StoragePathInfo}
@@ -28,9 +28,12 @@ import org.apache.spark.sql.SparkSession
 import org.apache.spark.sql.catalyst.InternalRow
 import org.apache.spark.sql.execution.PartitionedFileUtil
 import org.apache.spark.sql.execution.datasources.FilePartition
+import org.junit.jupiter.api.Test
 import org.junit.jupiter.params.ParameterizedTest
 import org.junit.jupiter.params.provider.ValueSource
 
+import scala.collection.JavaConverters._
+
 class TestPartitionDirectoryConverter extends SparkAdapterSupport {
 
   val blockSize = 1024
@@ -49,6 +52,7 @@ class TestPartitionDirectoryConverter extends 
SparkAdapterSupport {
     val options = Map(
       
s"${HoodieStorageConfig.LOGFILE_TO_PARQUET_COMPRESSION_RATIO_FRACTION.key()}" 
-> logFraction.toString
     )
+    val config = new HoodieConfig(TypedProperties.fromMap(options.asJava))
     // There are 4 cases for file slices:
     // 1. base file only
     // 2. log files only
@@ -81,7 +85,7 @@ class TestPartitionDirectoryConverter extends 
SparkAdapterSupport {
     val partitionValues = Seq("2025-01-01")
 
     val partitionedFiles = slices.flatMap(slice => {
-      val dir = 
PartitionDirectoryConverter.convertFileSliceToPartitionDirectory(InternalRow.fromSeq(partitionValues),
 slice, options)
+      val dir = 
PartitionDirectoryConverter.convertFileSliceToPartitionDirectory(InternalRow.fromSeq(partitionValues),
 slice, config)
       sparkAdapter.splitFiles(spark, dir, false, maxSplitSize)
     })
 
@@ -90,6 +94,25 @@ class TestPartitionDirectoryConverter extends 
SparkAdapterSupport {
     spark.stop()
   }
 
+  @Test
+  def testConvertNativeLogFileSliceToPartitionDirectory(): Unit = {
+    val logFraction = 0.1
+    val options = Map(
+      
s"${HoodieStorageConfig.LOGFILE_TO_PARQUET_COMPRESSION_RATIO_FRACTION.key()}" 
-> logFraction.toString
+    )
+    val config = new HoodieConfig(TypedProperties.fromMap(options.asJava))
+    val fileId = "native-log-file"
+    val recordCount = 300
+    val nativeLogFile = buildNativeHoodieLogFile(fileId, recordCount)
+    val slice = new FileSlice(new HoodieFileGroupId(partitionPath, fileId), 
baseInstant)
+    slice.addLogFile(nativeLogFile)
+
+    assert(nativeLogFile.isNativeLogFile)
+    val directory = 
PartitionDirectoryConverter.convertFileSliceToPartitionDirectory(
+      InternalRow.fromSeq(Seq("2025-01-01")), slice, config)
+    assert(directory.files.head.getLen == recordCount * 
fixedSizePerRecordWithParquetFormat)
+  }
+
   private def verifyBalanceByNum(tasks: Seq[FilePartition], totalRecordNum: 
Int, logFraction: Double): Unit = {
     val expectedRecordNumPerTask = totalRecordNum / tasks.size
     val expectedToleranceMin = expectedRecordNumPerTask * 0.9
@@ -161,4 +184,11 @@ class TestPartitionDirectoryConverter extends 
SparkAdapterSupport {
     new HoodieLogFile(info)
   }
 
+  private def buildNativeHoodieLogFile(fileId: String, recordsNum: Int): 
HoodieLogFile = {
+    val path = new StoragePath(s"${fileId}_0-0-0_${baseInstant}_1.log.parquet")
+    val fileLen = recordsNum * fixedSizePerRecordWithParquetFormat
+    val info = new StoragePathInfo(path, fileLen, false, 1, blockSize, 
System.currentTimeMillis())
+    new HoodieLogFile(info)
+  }
+
 }

Reply via email to