Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,6 @@ public abstract class BaseHoodieTableFileIndex implements AutoCloseable {
@Getter(AccessLevel.PROTECTED)
private final String[] partitionColumns;

@Getter
protected final HoodieMetadataConfig metadataConfig;
Comment thread
voonhous marked this conversation as resolved.
private final TypedProperties configProperties;
@Getter
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@

package org.apache.hudi.core.read;

import org.apache.hudi.common.config.HoodieMetadataConfig;
import org.apache.hudi.common.model.FileSlice;
import org.apache.hudi.core.read.BaseHoodieTableFileIndex.PartitionPath;
import org.apache.hudi.storage.StoragePath;
Expand All @@ -35,40 +34,11 @@

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.mock;

public class BaseHoodieTableFileIndexTest {

@Test
public void testGetMetadataConfigReturnsFieldValue() throws Exception {
// Create a mock of BaseHoodieTableFileIndex
BaseHoodieTableFileIndex fileIndex = mock(BaseHoodieTableFileIndex.class,
org.mockito.Mockito.CALLS_REAL_METHODS);

// Create a test metadata config
HoodieMetadataConfig testConfig = HoodieMetadataConfig.newBuilder()
.enable(true)
.withMetadataIndexBloomFilter(true)
.withMetadataIndexColumnStats(true)
.build();

// Use reflection to set the private metadataConfig field
Field metadataConfigField = BaseHoodieTableFileIndex.class.getDeclaredField("metadataConfig");
metadataConfigField.setAccessible(true);
metadataConfigField.set(fileIndex, testConfig);

// Test the getMetadataConfig method
HoodieMetadataConfig result = fileIndex.getMetadataConfig();

assertNotNull(result, "Metadata config should not be null");
assertSame(testConfig, result, "Should return the same metadata config instance");
assertEquals(true, result.isEnabled(), "Metadata should be enabled");
assertEquals(true, result.isBloomFilterIndexEnabled(), "Bloom filter index should be enabled");
assertEquals(true, result.isColumnStatsIndexEnabled(), "Column stats index should be enabled");
}

/**
* Regression test for the empty-partition NPE that surfaces in {@code getInputFileSlices}
* when the {@code hoodie.datasource.read.file.index.list.file.statuses.using.ro.path.filter}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,19 +19,23 @@
package org.apache.hudi.common.bootstrap.index;

import org.apache.hudi.common.config.HoodieCommonConfig;
import org.apache.hudi.common.config.HoodieMetadataConfig;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.engine.HoodieEngineContext;
import org.apache.hudi.common.engine.HoodieLocalEngineContext;
import org.apache.hudi.common.model.HoodieTableQueryType;
import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.table.HoodieTableMetaClient;
import org.apache.hudi.common.testutils.HoodieCommonTestHarness;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.core.read.BaseHoodieTableFileIndex;
import org.apache.hudi.metadata.HoodieTableMetadataUtil;
import org.apache.hudi.storage.StoragePath;

import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.CsvSource;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.junit.jupiter.MockitoExtension;

Expand Down Expand Up @@ -60,6 +64,7 @@ void testGetFileSlicesCount(boolean useSpillableMap) throws IOException {
Collections.emptyList(),
Option.empty(),
false,
false,
true,
null,
true,
Expand All @@ -69,13 +74,66 @@ void testGetFileSlicesCount(boolean useSpillableMap) throws IOException {
Assertions.assertTrue(baseHoodieTableFileIndex.getFileSlicesCount() == 0);
}

/**
* The metadata config the index builds is gated on a three-way conjunction, and both of the
* added conjuncts came from behavior fixes: {@code isFilesPartitionAvailable} from HUDI-5403
* (#7488, a Trino listing regression) and {@code useLatestBaseFilesPathFilterForListing} from
* #18136. Nothing else in the repo asserts how that conjunction resolves.
*/
@ParameterizedTest
@CsvSource({
"true, true, false, true",
"true, true, true, false",
"true, false, false, false",
"true, false, true, false",
"false, true, false, false",
"false, true, true, false",
"false, false, false, false",
"false, false, true, false"
})
void testMetadataConfigIsEnabledOnlyWhenEveryConjunctHolds(boolean metadataEnabled,
boolean filesPartitionAvailable,
boolean useLatestBaseFilesPathFilterForListing,
boolean expectedEnabled) throws IOException {
initMetaClient();
if (filesPartitionAvailable) {
metaClient.getTableConfig().setValue(
HoodieTableConfig.TABLE_METADATA_PARTITIONS, HoodieTableMetadataUtil.PARTITION_NAME_FILES);
}
TypedProperties properties = new TypedProperties();
properties.put(HoodieMetadataConfig.ENABLE.key(), String.valueOf(metadataEnabled));

TestLocalIndex fileIndex = new TestLocalIndex(
new HoodieLocalEngineContext(getStorageConf()),
metaClient,
properties,
HoodieTableQueryType.READ_OPTIMIZED,
Collections.emptyList(),
Option.empty(),
useLatestBaseFilesPathFilterForListing,
false,
true,
null,
true,
Option.empty(),
Option.empty()
);

Assertions.assertEquals(expectedEnabled, fileIndex.metadataConfig().isEnabled());
}

private static class TestLocalIndex extends BaseHoodieTableFileIndex {

public TestLocalIndex(HoodieEngineContext engineContext, HoodieTableMetaClient metaClient, TypedProperties configProperties, HoodieTableQueryType queryType,
List<StoragePath> queryPaths, Option<String> specifiedQueryInstant, boolean shouldIncludePendingCommits, boolean shouldValidateInstant,
List<StoragePath> queryPaths, Option<String> specifiedQueryInstant, boolean useLatestBaseFilesPathFilterForListing,
boolean shouldIncludePendingCommits, boolean shouldValidateInstant,
FileStatusCache fileStatusCache, boolean shouldListLazily, Option<String> startCompletionTime, Option<String> endCompletionTime) {
super(engineContext, metaClient, configProperties, queryType, queryPaths, specifiedQueryInstant, false, shouldIncludePendingCommits, shouldValidateInstant, fileStatusCache, shouldListLazily,
startCompletionTime, endCompletionTime);
super(engineContext, metaClient, configProperties, queryType, queryPaths, specifiedQueryInstant, useLatestBaseFilesPathFilterForListing, shouldIncludePendingCommits, shouldValidateInstant,
fileStatusCache, shouldListLazily, startCompletionTime, endCompletionTime);
}

HoodieMetadataConfig metadataConfig() {
return metadataConfig;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ import org.apache.hudi.HoodieBaseRelation.{convertToHoodieSchema, createHFileRea
import org.apache.hudi.HoodieConversionUtils.toScalaOption
import org.apache.hudi.client.utils.SparkInternalSchemaConverter
import org.apache.hudi.common.avro.HoodieAvroUtils
import org.apache.hudi.common.config.{ConfigProperty, HoodieMetadataConfig}
import org.apache.hudi.common.config.ConfigProperty
import org.apache.hudi.common.fs.FSUtils
import org.apache.hudi.common.fs.FSUtils.getRelativePartitionPath
import org.apache.hudi.common.model.{FileSlice, HoodieFileFormat, HoodieRecord}
Expand All @@ -40,7 +40,6 @@ import org.apache.hudi.common.util.HoodieStorageUtils
import org.apache.hudi.common.util.StringUtils.isNullOrEmpty
import org.apache.hudi.common.util.ValidationUtils.checkState
import org.apache.hudi.config.HoodieBootstrapConfig.DATA_QUERIES_ONLY
import org.apache.hudi.config.HoodieWriteConfig
import org.apache.hudi.exception.HoodieException
import org.apache.hudi.hadoop.fs.HadoopFSUtils
import org.apache.hudi.hadoop.fs.HadoopFSUtils.convertToStoragePath
Expand Down Expand Up @@ -75,16 +74,6 @@ trait HoodieFileSplit {}

case class HoodieTableSchema(structTypeSchema: StructType, schema: HoodieSchema, internalSchema: Option[InternalSchema] = None)

case class HoodieTableState(tablePath: String,
latestCommitTimestamp: Option[String],
recordKeyField: String,
orderingFields: List[String],
usesVirtualKeys: Boolean,
metadataConfig: HoodieMetadataConfig,
recordMergeImplClasses: List[String],
recordMergeStrategyId: String)
Comment thread
ryux1 marked this conversation as resolved.


/**
* Hoodie BaseRelation which extends [[PrunedFilteredScan]]
*/
Expand Down Expand Up @@ -244,28 +233,16 @@ abstract class HoodieBaseRelation(val sqlContext: SQLContext,
/**
* NOTE: PLEASE READ THIS CAREFULLY
*
* Even though [[HoodieFileIndex]] initializes eagerly listing all of the files w/in the given Hudi table,
* this variable itself is _lazy_ (and have to stay that way) which guarantees that it's not initialized, until
* Even though [[HoodieFileIndex]] does eager work on construction (it opens the metadata-table reader
* and reloads the active timeline, and lists every file in the table when
* `hoodie.datasource.read.file.index.listing.mode` is `eager`; it defaults to `lazy`), this variable
* itself is _lazy_ (and have to stay that way) which guarantees that it's not initialized, until
* it's actually accessed
*/
lazy val fileIndex: HoodieFileIndex =
HoodieFileIndex(sparkSession, metaClient, Some(tableStructSchema), optParams,
FileStatusCache.getOrCreate(sparkSession), shouldIncludeLogFiles())
Comment thread
voonhous marked this conversation as resolved.

lazy val tableState: HoodieTableState = {
val recordMergerImpls = optParams.get(HoodieWriteConfig.RECORD_MERGE_IMPL_CLASSES.key()).map(impls => ConfigUtils.split2List(impls).asScala.toList).getOrElse(List.empty)
// Subset of the state of table's configuration as of at the time of the query
HoodieTableState(tablePath = basePath.toString,
latestCommitTimestamp = queryTimestamp,
recordKeyField = recordKeyField,
orderingFields = orderingFields,
usesVirtualKeys = !tableConfig.populateMetaFields(),
metadataConfig = fileIndex.getMetadataConfig,
Comment thread
ryux1 marked this conversation as resolved.
Comment thread
voonhous marked this conversation as resolved.
recordMergeImplClasses = recordMergerImpls,
recordMergeStrategyId = tableConfig.getRecordMergeStrategyId
)
}

/**
* Columns that relation has to read from the storage to properly execute on its semantic: for ex,
* for Merge-on-Read tables key fields as well and precombine field comprise mandatory set of columns,
Expand All @@ -280,7 +257,7 @@ abstract class HoodieBaseRelation(val sqlContext: SQLContext,
// NOTE: We're including compaction here since it's not considering a "commit" operation
metaClient.getCommitsAndCompactionTimeline.filterCompletedInstants

private def queryTimestamp: Option[String] =
protected def queryTimestamp: Option[String] =
specifiedQueryTimestamp.orElse(toScalaOption(timeline.lastInstant()).map(_.requestedTime))

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -86,12 +86,16 @@ private[hudi] case class HoodieMergeOnReadBaseFileReaders(fullSchemaReader: Base
*
* @param sc spark's context
* @param config hadoop configuration
* @param sqlConf spark SQL configuration
* @param fileReaders suite of base file readers
* @param tableSchema table's full schema
* @param requiredSchema expected (potentially) projected schema
* @param tableState table's state
* @param targetInstantTime as-of instant, or last instant in the query timeline
* @param mergeType type of merge performed
* @param fileSplits target file-splits this RDD will be iterating over
* @param optionalFilters filters to apply while reading records
* @param metaClient table metadata client
* @param options datasource options
* @param includedInstantTimeSet instant time set used to filter records
*/
class HoodieMergeOnReadRDDV2(@transient sc: SparkContext,
Expand All @@ -100,7 +104,7 @@ class HoodieMergeOnReadRDDV2(@transient sc: SparkContext,
fileReaders: HoodieMergeOnReadBaseFileReaders,
tableSchema: HoodieTableSchema,
requiredSchema: HoodieTableSchema,
tableState: HoodieTableState,
targetInstantTime: Option[String],
mergeType: String,
@transient fileSplits: Seq[HoodieMergeOnReadFileSplit],
optionalFilters: Array[Filter],
Expand Down Expand Up @@ -208,7 +212,7 @@ class HoodieMergeOnReadRDDV2(@transient sc: SparkContext,
val fileGroupReader: HoodieFileGroupReader[IndexedRecord] = HoodieFileGroupReader.builder()
.withReaderContext(readerContext)
.withHoodieTableMetaClient(metaClient)
.withLatestCommitTime(tableState.latestCommitTimestamp.orNull)
.withLatestCommitTime(targetInstantTime.orNull)
.withLogFiles(logFiles.stream())
.withBaseFileOption(baseFileOption)
.withPartitionPath(partitionPath)
Expand All @@ -226,7 +230,7 @@ class HoodieMergeOnReadRDDV2(@transient sc: SparkContext,
HoodieLsmFileGroupReader.builder[InternalRow]()
.withReaderContext(readerContext)
.withHoodieTableMetaClient(metaClient)
.withLatestCommitTime(tableState.latestCommitTimestamp.orNull)
.withLatestCommitTime(targetInstantTime.orNull)
.withLogFiles(logFiles.stream())
.withBaseFileOption(baseFileOption)
.withPartitionPath(partitionPath)
Expand All @@ -239,7 +243,7 @@ class HoodieMergeOnReadRDDV2(@transient sc: SparkContext,
HoodieFileGroupReader.builder[InternalRow]()
.withReaderContext(readerContext)
.withHoodieTableMetaClient(metaClient)
.withLatestCommitTime(tableState.latestCommitTimestamp.orNull)
.withLatestCommitTime(targetInstantTime.orNull)
.withLogFiles(logFiles.stream())
.withBaseFileOption(baseFileOption)
.withPartitionPath(partitionPath)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,7 @@ case class MergeOnReadIncrementalRelationV1(override val sqlContext: SQLContext,
fileReaders = readers,
tableSchema = tableSchema,
requiredSchema = requiredSchema,
tableState = tableState,
targetInstantTime = targetInstantTime,
mergeType = mergeType,
fileSplits = fileSplits,
includedInstantTimeSet = Option(includedCommits.map(_.requestedTime).toSet),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,7 @@ case class MergeOnReadIncrementalRelationV2(override val sqlContext: SQLContext,
fileReaders = readers,
tableSchema = tableSchema,
requiredSchema = requiredSchema,
tableState = tableState,
targetInstantTime = targetInstantTime,
mergeType = mergeType,
fileSplits = fileSplits,
includedInstantTimeSet = Option(includedCommits.map(_.requestedTime).toSet),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,16 @@ abstract class BaseMergeOnReadSnapshotRelation(sqlContext: SQLContext,
protected val mergeType: String = optParams.getOrElse(DataSourceReadOptions.REALTIME_MERGE.key,
DataSourceReadOptions.REALTIME_MERGE.defaultValue)

/**
* Instant this query is targeting: the as-of instant when time-travel is requested, otherwise
* the last instant of the (potentially narrowed) query timeline.
*
* NOTE: This is deliberately a `lazy val` rather than a `def`, so the instant is captured once
* per relation instance instead of being re-derived from [[timeline]] on each
* [[composeRDD]]. This preserves the behavior of the `lazy val tableState` it replaces.
*/
protected lazy val targetInstantTime: Option[String] = queryTimestamp

protected override def composeRDD(fileSplits: Seq[HoodieMergeOnReadFileSplit],
tableSchema: HoodieTableSchema,
requiredSchema: HoodieTableSchema,
Expand All @@ -108,7 +118,7 @@ abstract class BaseMergeOnReadSnapshotRelation(sqlContext: SQLContext,
fileReaders = readers,
tableSchema = tableSchema,
requiredSchema = requiredSchema,
tableState = tableState,
targetInstantTime = targetInstantTime,
mergeType = mergeType,
fileSplits = fileSplits,
optionalFilters = optionalFilters,
Expand Down
Loading