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.");
+ }
+ }
}