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

yihua 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 34b4850421eb fix(streamer): route configured write table version into 
sample-writes flow (#19746)
34b4850421eb is described below

commit 34b4850421eb4c27d63679204e0014137bb1ab20
Author: Lokesh Jain <[email protected]>
AuthorDate: Thu Aug 27 04:00:14 2026 +0530

    fix(streamer): route configured write table version into sample-writes flow 
(#19746)
---
 .../utilities/streamer/SparkSampleWritesUtils.java |  5 ++++
 .../deltastreamer/TestSparkSampleWritesUtils.java  | 33 ++++++++++++++++++++--
 2 files changed, 36 insertions(+), 2 deletions(-)

diff --git 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/SparkSampleWritesUtils.java
 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/SparkSampleWritesUtils.java
index b85d5241196d..e5c216305b9d 100644
--- 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/SparkSampleWritesUtils.java
+++ 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/SparkSampleWritesUtils.java
@@ -91,9 +91,13 @@ public class SparkSampleWritesUtils {
       throws IOException {
     String uniqueId = UUID.randomUUID().toString();
     final String sampleWritesBasePath = getSampleWritesBasePath(jsc, 
writeConfig, uniqueId);
+    // Propagate the user's configured write table version to the 
sample-writes shadow table so its
+    // on-disk layout matches the version the inherited write config (and the 
SparkRDDWriteClient
+    // below) operates with.
     HoodieTableMetaClient.newTableBuilder()
         .setTableType(HoodieTableType.COPY_ON_WRITE)
         .setTableName(String.format("%s_samples_%s", 
writeConfig.getTableName(), uniqueId))
+        .setTableVersion(writeConfig.getWriteVersion())
         .setCDCEnabled(false)
         
.initTable(HadoopFSUtils.getStorageConfWithCopy(jsc.hadoopConfiguration()), 
sampleWritesBasePath);
     TypedProperties props = writeConfig.getProps();
@@ -105,6 +109,7 @@ public class SparkSampleWritesUtils {
         .withSchemaEvolutionEnable(false)
         .withBulkInsertParallelism(1)
         .withPath(sampleWritesBasePath)
+        .withWriteTableVersion(writeConfig.getWriteVersion().versionCode())
         .build();
     Pair<Boolean, String> emptyRes = Pair.of(false, null);
     try (SparkRDDWriteClient sampleWriteClient = new SparkRDDWriteClient(new 
HoodieSparkEngineContext(jsc), sampleWriteConfig, Option.empty())) {
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestSparkSampleWritesUtils.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestSparkSampleWritesUtils.java
index 33d084903e46..7d62716a7a20 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestSparkSampleWritesUtils.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestSparkSampleWritesUtils.java
@@ -23,11 +23,13 @@ import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieTableType;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.HoodieTableVersion;
 import org.apache.hudi.common.testutils.HoodieTestDataGenerator;
 import org.apache.hudi.common.testutils.HoodieTestTable;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.config.HoodieCompactionConfig;
 import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.hadoop.fs.HadoopFSUtils;
 import org.apache.hudi.testutils.SparkClientFunctionalTestHarness;
 import org.apache.hudi.utilities.config.HoodieStreamerConfig;
 import org.apache.hudi.utilities.streamer.SparkSampleWritesUtils;
@@ -39,6 +41,8 @@ import org.apache.spark.api.java.JavaRDD;
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
 
 import java.io.IOException;
 import java.util.ArrayList;
@@ -92,17 +96,20 @@ public class TestSparkSampleWritesUtils extends 
SparkClientFunctionalTestHarness
     assertEquals(originalRecordSize, 
originalWriteConfig.getCopyOnWriteRecordSizeEstimate(), "Original record size 
estimate should not be changed.");
   }
 
-  @Test
-  void overwriteRecordSizeEstimateForEmptyTable() throws IOException {
+  @ParameterizedTest
+  @EnumSource(value = HoodieTableVersion.class, names = {"SIX", "NINE"})
+  void overwriteRecordSizeEstimateForEmptyTable(HoodieTableVersion 
tableVersion) throws IOException {
     int originalRecordSize = 100;
     TypedProperties props = new TypedProperties();
     props.put(HoodieStreamerConfig.SAMPLE_WRITES_ENABLED.key(), "true");
     props.put(HoodieCompactionConfig.COPY_ON_WRITE_RECORD_SIZE_ESTIMATE.key(), 
String.valueOf(originalRecordSize));
+    props.put(HoodieWriteConfig.WRITE_TABLE_VERSION.key(), 
String.valueOf(tableVersion.versionCode()));
     HoodieWriteConfig originalWriteConfig = HoodieWriteConfig.newBuilder()
         .withProperties(props)
         .forTable("foo")
         .withPath(basePath())
         .withSchema(HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA)
+        .withWriteTableVersion(tableVersion.versionCode())
         .build();
 
     String commitTime = HoodieTestDataGenerator.getCommitTimeAtUTC(1);
@@ -112,6 +119,7 @@ public class TestSparkSampleWritesUtils extends 
SparkClientFunctionalTestHarness
     assertTrue(writeConfigOpt.isPresent());
     assertEquals(337.0, 
writeConfigOpt.get().getCopyOnWriteRecordSizeEstimate(), 10.0);
     assertSampleWritesNonPartitioned();
+    assertSampleWritesShadowTableVersion(tableVersion);
   }
 
   @Test
@@ -165,4 +173,25 @@ public class TestSparkSampleWritesUtils extends 
SparkClientFunctionalTestHarness
               + partitionDirs);
     }
   }
+
+  /**
+   * Fails if any sample-writes shadow table on disk was not created at the 
expected table version,
+   * i.e. verifies the configured write version was routed into the shadow 
table instead of
+   * defaulting to the current version.
+   */
+  private void assertSampleWritesShadowTableVersion(HoodieTableVersion 
expected) throws IOException {
+    Path sampleWritesPath = new Path(basePath(), SAMPLE_WRITES_FOLDER_PATH);
+    FileSystem fs = 
sampleWritesPath.getFileSystem(jsc().hadoopConfiguration());
+    assertTrue(fs.exists(sampleWritesPath), "Sample-writes folder should exist 
after a sample write.");
+    FileStatus[] runs = fs.listStatus(sampleWritesPath);
+    assertTrue(runs.length > 0, "Sample-writes folder should contain at least 
one run.");
+    for (FileStatus run : runs) {
+      HoodieTableMetaClient sampleMetaClient = HoodieTableMetaClient.builder()
+          
.setConf(HadoopFSUtils.getStorageConfWithCopy(jsc().hadoopConfiguration()))
+          .setBasePath(run.getPath().toString())
+          .build();
+      assertEquals(expected, 
sampleMetaClient.getTableConfig().getTableVersion(),
+          "Sample-writes shadow table at " + run.getPath() + " should be 
created at the configured write version.");
+    }
+  }
 }

Reply via email to