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 133a25515676 fix(utilities): include metadata table index-init
instants in the record index validation snapshot (#19395)
133a25515676 is described below
commit 133a25515676ba31a2a258147dc8d563fb280de1
Author: Davis-Zhang-Onehouse
<[email protected]>
AuthorDate: Fri Jul 31 02:25:39 2026 -0700
fix(utilities): include metadata table index-init instants in the record
index validation snapshot (#19395)
HoodieMetadataTableValidator reads the metadata table with a time-travel
snapshot anchored to the data table's latest completed commit. On a table
version 6 metadata table the partition-initialization deltacommits are that
same data instant with a three-digit suffix appended (010 for FILES, 011 for
RECORD_INDEX), and Hudi compares instants as strings, so those derived
instants sort after the bare data instant and fall outside the snapshot.
When the data table has no commit newer than the initialization instant,
every record index file slice is filtered out, the index reads back empty,
and the validator reports 100% of the data table's keys as missing from it.
This is permanent for a table that has stopped receiving writes; cleans and
rollbacks do not help because getWriteTimeline() only whitelists commit,
deltacommit, compaction, logcompaction and replacecommit.
Table version 8 and above are unaffected: generateUniqueInstantTime derives
the init instants from SOLO_COMMIT_TIMESTAMP, which sorts below every data
instant.
Advance the instant used for the metadata table read by one millisecond.
That is strictly greater than any <instant><suffix>, which shares the whole
17-character prefix, while remaining a valid yyyyMMddHHmmssSSS instant - the
time travel option runs the value through formatQueryInstant, which rejects
anything else. The data table side is unchanged, so both sides stay pinned
to
the same data instant.
Fixes both validateRecordIndexContent and validateRecordIndexCount; the
latter carried the same defect, masked only because the content check
shadows it.
---
.../utilities/HoodieMetadataTableValidator.java | 42 ++++++-
.../TestHoodieMetadataTableValidator.java | 122 +++++++++++++++++++++
2 files changed, 162 insertions(+), 2 deletions(-)
diff --git
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieMetadataTableValidator.java
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieMetadataTableValidator.java
index 3f28b4f14187..fa99795c22fc 100644
---
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieMetadataTableValidator.java
+++
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/HoodieMetadataTableValidator.java
@@ -61,6 +61,7 @@ import
org.apache.hudi.common.table.timeline.HoodieActiveTimeline;
import org.apache.hudi.common.table.timeline.HoodieInstant;
import org.apache.hudi.common.table.timeline.HoodieTimeline;
import org.apache.hudi.common.table.timeline.InstantComparison;
+import org.apache.hudi.common.table.timeline.TimelineUtils;
import org.apache.hudi.common.table.view.FileSystemViewManager;
import org.apache.hudi.common.table.view.FileSystemViewStorageConfig;
import org.apache.hudi.common.table.view.FileSystemViewStorageType;
@@ -111,11 +112,13 @@ import org.apache.spark.storage.StorageLevel;
import java.io.IOException;
import java.io.Serializable;
import java.nio.ByteBuffer;
+import java.text.ParseException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
import java.util.Comparator;
+import java.util.Date;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@@ -208,6 +211,10 @@ public class HoodieMetadataTableValidator implements
Serializable {
private static final long serialVersionUID = 1L;
+ // Advance the metadata table query instant by this much so that instants
derived from a data table
+ // instant, which carry a three-digit suffix, fall inside the queried
window. See #metadataTableInstantFor.
+ private static final long METADATA_INSTANT_LOOKAHEAD_MS = 1;
+
// Spark context
private transient JavaSparkContext jsc;
// config
@@ -1235,7 +1242,7 @@ public class HoodieMetadataTableValidator implements
Serializable {
.select(RECORD_KEY_METADATA_FIELD)
.count();
long countKeyFromRecordIndex =
sparkEngineContext.getSqlContext().read().format("hudi")
-
.option(DataSourceReadOptions.TIME_TRAVEL_AS_OF_INSTANT().key(),latestCompletedCommit)
+ .option(DataSourceReadOptions.TIME_TRAVEL_AS_OF_INSTANT().key(),
metadataTableInstantFor(latestCompletedCommit))
.load(getMetadataTableBasePath(basePath))
.select("key")
.filter("type = 5")
@@ -1337,6 +1344,37 @@ public class HoodieMetadataTableValidator implements
Serializable {
}
}
+ /**
+ * Returns the instant to query the metadata table with, so that the
snapshot reflects the data
+ * table as of {@code dataTableInstant}.
+ * <p>
+ * Metadata table instants derived from a data table instant carry a
three-digit numeric suffix:
+ * partition initialization appends 010 and up (see
+ * {@code HoodieTableMetadataUtil#createIndexInitTimestamp}), and
metadata-table-internal
+ * compaction, clean, restore, indexing, log compaction and rollback append
001 to 006. Hudi
+ * compares instants as strings, so every one of those derived instants
sorts AFTER the bare data
+ * instant, and a snapshot taken as of the data instant itself excludes them
- leaving, for
+ * instance, the record index unreadable until the data table receives
another commit.
+ * <p>
+ * The bound is therefore advanced by a single millisecond. That is strictly
greater than any
+ * {@code <dataTableInstant><suffix>} (which shares the whole 17-character
prefix and so compares
+ * lower), while still being a valid {@code yyyyMMddHHmmssSSS} instant - the
time travel option
+ * rejects anything else, see {@code
HoodieSqlCommonUtils#formatQueryInstant}. Instants that are
+ * not timestamps (legacy or test instants such as "100") are returned
unchanged; they have no
+ * metadata table counterpart to include.
+ */
+ @VisibleForTesting
+ static String metadataTableInstantFor(String dataTableInstant) {
+ try {
+ Date dataTableInstantDate =
TimelineUtils.parseDateFromInstantTime(dataTableInstant);
+ return TimelineUtils.formatDate(new Date(dataTableInstantDate.getTime()
+ METADATA_INSTANT_LOOKAHEAD_MS));
+ } catch (ParseException e) {
+ log.warn("Cannot parse instant {} as a timestamp; querying the metadata
table as of it verbatim",
+ dataTableInstant);
+ return dataTableInstant;
+ }
+ }
+
@VisibleForTesting
JavaPairRDD<String, Pair<String, String>>
getRecordLocationsFromFSBasedListing(HoodieSparkEngineContext
sparkEngineContext,
String basePath,
@@ -1357,7 +1395,7 @@ public class HoodieMetadataTableValidator implements
Serializable {
String
basePath,
String
latestCompletedCommit) {
return sparkEngineContext.getSqlContext().read().format("hudi")
- .option(DataSourceReadOptions.TIME_TRAVEL_AS_OF_INSTANT().key(),
latestCompletedCommit)
+ .option(DataSourceReadOptions.TIME_TRAVEL_AS_OF_INSTANT().key(),
metadataTableInstantFor(latestCompletedCommit))
.load(getMetadataTableBasePath(basePath))
.filter("type = 5")
.select(functions.col("key"),
diff --git
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/TestHoodieMetadataTableValidator.java
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/TestHoodieMetadataTableValidator.java
index f16b0f49fb12..f3eaa254eeb5 100644
---
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/TestHoodieMetadataTableValidator.java
+++
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/TestHoodieMetadataTableValidator.java
@@ -37,6 +37,7 @@ import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.WriteOperationType;
import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.HoodieTableVersion;
import org.apache.hudi.common.table.log.HoodieLogFormat;
import org.apache.hudi.common.table.log.HoodieLogFormatWriter;
import org.apache.hudi.common.table.log.block.HoodieAvroDataBlock;
@@ -45,6 +46,7 @@ import org.apache.hudi.common.table.log.block.HoodieLogBlock;
import org.apache.hudi.common.table.timeline.HoodieActiveTimeline;
import org.apache.hudi.common.table.timeline.HoodieInstant;
import org.apache.hudi.common.table.timeline.HoodieTimeline;
+import org.apache.hudi.common.table.timeline.InstantComparison;
import org.apache.hudi.common.table.timeline.TimeGenerator;
import org.apache.hudi.common.table.timeline.TimeGenerators;
import org.apache.hudi.common.table.timeline.TimelineUtils;
@@ -57,9 +59,14 @@ import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.StringUtils;
import org.apache.hudi.common.util.collection.Pair;
import org.apache.hudi.config.HoodieCompactionConfig;
+import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.exception.HoodieIOException;
import org.apache.hudi.exception.HoodieValidationException;
import org.apache.hudi.hadoop.fs.HadoopFSUtils;
+import org.apache.hudi.metadata.HoodieTableMetadata;
+import org.apache.hudi.metadata.HoodieTableMetadataWriter;
+import org.apache.hudi.metadata.MetadataPartitionType;
+import org.apache.hudi.metadata.SparkMetadataWriterFactory;
import org.apache.hudi.metadata.stats.HoodieColumnRangeMetadata;
import org.apache.hudi.metadata.stats.ValueMetadata;
import org.apache.hudi.storage.HoodieStorage;
@@ -1678,4 +1685,119 @@ public class TestHoodieMetadataTableValidator extends
HoodieSparkClientTestBase
when(storage.listFiles(new StoragePath(basePath + "/" +
partition))).thenReturn(Collections.singletonList(storagePathInfo));
}
}
+
+ /**
+ * On a table version 6 metadata table the partition-initialization
deltacommits are the data instant
+ * they were derived from with a three-digit suffix appended (010 for FILES,
011 for RECORD_INDEX) -
+ * see {@code
HoodieBackedTableMetadataWriterTableVersionSix#createIndexInitTimestamp}. Those
instants
+ * sort AFTER the bare data instant under Hudi's lexicographic instant
comparison, so a metadata-table
+ * snapshot taken as of the data table's latest completed commit must still
include them. When the data
+ * table has no commit after the metadata table was initialized, failing to
do so makes the record index
+ * read back empty and every data-table key is reported as missing from it.
+ * <p>
+ * Table version 8 and above are unaffected: {@code
HoodieBackedTableMetadataWriter#generateUniqueInstantTime}
+ * derives the init instants from {@code SOLO_COMMIT_TIMESTAMP}, which sorts
below every data instant.
+ */
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ public void
testRecordIndexValidationWhenMdtInitializedAtLatestDataCommit(boolean
validateContent) throws Exception {
+ Map<String, String> writeOptions = new HashMap<>();
+ writeOptions.put(DataSourceWriteOptions.TABLE_NAME().key(), "test_table");
+ writeOptions.put("hoodie.table.name", "test_table");
+ writeOptions.put(DataSourceWriteOptions.TABLE_TYPE().key(),
"COPY_ON_WRITE");
+ writeOptions.put(DataSourceWriteOptions.RECORDKEY_FIELD().key(),
"_row_key");
+ writeOptions.put(DataSourceWriteOptions.PRECOMBINE_FIELD().key(),
"timestamp");
+ writeOptions.put(DataSourceWriteOptions.OPERATION().key(),
WriteOperationType.BULK_INSERT.value());
+ // Leave the metadata table off for the write; it is bootstrapped out of
band below so that the
+ // data table's latest completed commit stays the instant the init
instants are derived from.
+ writeOptions.put(HoodieMetadataConfig.ENABLE.key(), "false");
+ // Table version 6 is the one whose metadata table initialization instants
carry the suffix.
+ writeOptions.put(HoodieWriteConfig.WRITE_TABLE_VERSION.key(), "6");
+
+ makeInsertDf("000", 50).write().format("hudi").options(writeOptions)
+ .mode(SaveMode.Overwrite)
+ .save(basePath);
+
+ HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+ .withPath(basePath)
+ .forTable("test_table")
+ .withWriteTableVersion(6)
+ .withMetadataConfig(HoodieMetadataConfig.newBuilder()
+ .enable(true)
+ .withEnableGlobalRecordLevelIndex(true)
+ .withRecordIndexFileGroupCount(1, 1)
+ .build())
+ .build();
+ HoodieTableMetaClient metaClientBeforeInit =
HoodieTableMetaClient.builder()
+
.setBasePath(basePath).setConf(HadoopFSUtils.getStorageConf(jsc.hadoopConfiguration())).build();
+ assertEquals(HoodieTableVersion.SIX,
metaClientBeforeInit.getTableConfig().getTableVersion(),
+ "the suffixed initialization instants only exist on table version 6");
+
+ // Creating the writer initializes the FILES and RECORD_INDEX partitions
from the filesystem
+ // without adding a commit to the data table. Go through the factory so
the table-version-6 writer
+ // is selected, exactly as production does.
+ try (HoodieTableMetadataWriter ignored = SparkMetadataWriterFactory.create(
+ HadoopFSUtils.getStorageConf(jsc.hadoopConfiguration()), writeConfig,
context,
+ Option.empty(), metaClientBeforeInit.getTableConfig())) {
+ // constructing the writer performs the initialization
+ }
+
+ HoodieTableMetaClient dataMetaClient =
HoodieTableMetaClient.reload(metaClientBeforeInit);
+
assertTrue(dataMetaClient.getTableConfig().isMetadataPartitionAvailable(MetadataPartitionType.RECORD_INDEX),
+ "record index should be registered on the data table");
+ String latestDataCommit =
dataMetaClient.getActiveTimeline().getCommitsAndCompactionTimeline()
+ .filterCompletedInstants().lastInstant().get().requestedTime();
+ HoodieTableMetaClient mdtMetaClient = HoodieTableMetaClient.builder()
+ .setBasePath(HoodieTableMetadata.getMetadataTableBasePath(basePath))
+
.setConf(HadoopFSUtils.getStorageConf(jsc.hadoopConfiguration())).build();
+ List<String> mdtInstants =
mdtMetaClient.getActiveTimeline().filterCompletedInstants()
+
.getInstantsAsStream().map(HoodieInstant::requestedTime).collect(Collectors.toList());
+ assertTrue(
+ mdtInstants.stream().allMatch(instant ->
instant.startsWith(latestDataCommit)
+ && instant.length() > latestDataCommit.length()),
+ "expected every metadata instant to be a suffixed extension of " +
latestDataCommit
+ + " but got " + mdtInstants);
+
+ HoodieMetadataTableValidator.Config config = new
HoodieMetadataTableValidator.Config();
+ config.basePath = "file:" + basePath;
+ config.validateLatestFileSlices = true;
+ // validateRecordIndexContent shadows validateRecordIndexCount, so
toggling it exercises both paths.
+ config.validateRecordIndexContent = validateContent;
+ config.validateRecordIndexCount = true;
+ config.ignoreFailed = true;
+
+ HoodieMetadataTableValidator validator = new
HoodieMetadataTableValidator(jsc, config);
+ // Assert on the record index validation directly:
doMetadataTableValidation() reports any
+ // non-HoodieValidationException as a successful run, so run() alone
cannot distinguish a genuine
+ // pass from the read blowing up.
+ assertDoesNotThrow(() -> validator.validateRecordIndex(new
HoodieSparkEngineContext(jsc), dataMetaClient),
+ "record index validation should pass against an intact record index");
+ assertTrue(validator.run(), "validation should succeed against an intact
record index");
+ assertFalse(validator.hasValidationFailure(), () -> "unexpected validation
failures: "
+ + validator.getThrowables());
+ }
+
+ @Test
+ public void testMetadataTableInstantForIncludesEveryDerivedInstant() {
+ String dataTableInstant = "20231012054834279";
+
+ String queryInstant =
HoodieMetadataTableValidator.metadataTableInstantFor(dataTableInstant);
+
+ assertEquals("20231012054834280", queryInstant, "the bound should advance
the instant by one millisecond");
+ // 010-013 are appended when a metadata table partition is initialized,
001-006 by the
+ // metadata-table-internal operations; all of them must fall inside the
queried window.
+ for (String suffix : new String[] {"001", "002", "003", "004", "005",
"006", "010", "011", "012", "013"}) {
+ assertTrue(
+ InstantComparison.compareTimestamps(dataTableInstant + suffix,
InstantComparison.LESSER_THAN_OR_EQUALS, queryInstant),
+ () -> "metadata instant " + dataTableInstant + suffix + " should be
included by " + queryInstant);
+ }
+ // ... while the next data table instant stays outside it.
+ assertTrue(InstantComparison.compareTimestamps("20231012054834281",
InstantComparison.GREATER_THAN, queryInstant),
+ "a later data table instant should not be included");
+ }
+
+ @Test
+ public void testMetadataTableInstantForLeavesNonTimestampInstantUnchanged() {
+ assertEquals("100",
HoodieMetadataTableValidator.metadataTableInstantFor("100"));
+ }
}