This is an automated email from the ASF dual-hosted git repository.
xiangfu0 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new 257a0e7df2c Reuse processed table state across segment loads (#19571)
257a0e7df2c is described below
commit 257a0e7df2cee4b26cf89c417eef13f71047f1c4
Author: Xiang Fu <[email protected]>
AuthorDate: Wed Sep 23 17:24:39 2026 -0700
Reuse processed table state across segment loads (#19571)
* Reuse normalized table schemas from the existing table manager cache
* Use cached table config and schema for ordinary segment loads
* Cache processed index loading settings for ordinary segment loads
* Share cached loading config on read-only segment load paths
* Derive segment-local index loading wrappers
* Separate immutable index loading state
* Separate resolved index state from segment overrides
* add regression test
* Fix default-tier reload config race
* Expose cached loading config to subclasses
* Avoid invalidating unchanged index state
---------
Co-authored-by: J-HowHuang <[email protected]>
---
.../core/data/manager/BaseTableDataManager.java | 22 +-
.../manager/offline/OfflineTableDataManager.java | 3 +-
.../manager/realtime/RealtimeTableDataManager.java | 12 +-
.../data/manager/BaseTableDataManagerTest.java | 152 +++++++
.../immutable/ImmutableSegmentLoader.java | 5 +-
.../segment/index/loader/IndexLoadingConfig.java | 438 +++++++++++++--------
.../index/loader/IndexLoadingConfigTest.java | 62 ++-
7 files changed, 511 insertions(+), 183 deletions(-)
diff --git
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
index 8a3269d6eb1..da5fa943ec3 100644
---
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
+++
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
@@ -176,6 +176,7 @@ public abstract class BaseTableDataManager implements
TableDataManager {
// Caches the latest TableConfig and Schema pair. The cache should not be
modified.
protected volatile Pair<TableConfig, Schema> _cachedTableConfigAndSchema;
+ protected volatile IndexLoadingConfig _cachedIndexLoadingConfig;
protected volatile boolean _shutDown;
protected volatile boolean _isDeleted;
@@ -437,6 +438,21 @@ public abstract class BaseTableDataManager implements
TableDataManager {
return indexLoadingConfig;
}
+ /// Returns the shared processed config for ordinary loads. Segment-specific
settings are applied with the
+ /// `with...` methods on [IndexLoadingConfig], which share its processed
table-level state.
+ /// Explicit reloads and config/schema refresh messages still use
[#fetchIndexLoadingConfig()].
+ protected IndexLoadingConfig getCachedIndexLoadingConfig() {
+ Pair<TableConfig, Schema> cached = _cachedTableConfigAndSchema;
+ IndexLoadingConfig indexLoadingConfig = _cachedIndexLoadingConfig;
+ if (indexLoadingConfig == null || indexLoadingConfig.getTableConfig() !=
cached.getLeft()
+ || indexLoadingConfig.getSchema() != cached.getRight()) {
+ indexLoadingConfig = new IndexLoadingConfig(_instanceDataManagerConfig,
cached.getLeft(), cached.getRight());
+ indexLoadingConfig.setTableDataDir(_tableDataDir);
+ _cachedIndexLoadingConfig = indexLoadingConfig;
+ }
+ return indexLoadingConfig;
+ }
+
@Override
public Pair<TableConfig, Schema> getCachedTableConfigAndSchema() {
return _cachedTableConfigAndSchema;
@@ -444,6 +460,10 @@ public abstract class BaseTableDataManager implements
TableDataManager {
@Override
public void updateCachedTableConfigAndSchema(TableConfig tableConfig, Schema
schema) {
+ // Normalize before publishing so segment loads never add timestamp fields
to a shared schema.
+ if (schema != null) {
+ TimestampIndexUtils.applyTimestampIndex(tableConfig, schema);
+ }
_cachedTableConfigAndSchema = Pair.of(tableConfig, schema);
}
@@ -1100,7 +1120,7 @@ public abstract class BaseTableDataManager implements
TableDataManager {
SegmentMetadata localMetadata, boolean forceDownload)
throws Exception {
String segmentTier = getSegmentCurrentTier(segmentName);
- indexLoadingConfig.setSegmentTier(segmentTier);
+ indexLoadingConfig = indexLoadingConfig.copyWithSegmentTier(segmentTier);
indexLoadingConfig.setTableDataDir(_tableDataDir);
File indexDir = getSegmentDataDir(segmentName, segmentTier,
indexLoadingConfig.getTableConfig());
_segmentReloadSemaphore.acquire(segmentName, _logger);
diff --git
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/OfflineTableDataManager.java
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/OfflineTableDataManager.java
index 8717408446f..74ed8c44e81 100644
---
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/OfflineTableDataManager.java
+++
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/OfflineTableDataManager.java
@@ -71,8 +71,7 @@ public class OfflineTableDataManager extends
BaseTableDataManager {
protected void doAddOnlineSegment(String segmentName)
throws Exception {
SegmentZKMetadata zkMetadata = fetchZKMetadata(segmentName);
- IndexLoadingConfig indexLoadingConfig = fetchIndexLoadingConfig();
- indexLoadingConfig.setSegmentTier(zkMetadata.getTier());
+ IndexLoadingConfig indexLoadingConfig =
getCachedIndexLoadingConfig().withSegmentTier(zkMetadata.getTier());
handleUpsertPreload(zkMetadata, indexLoadingConfig);
SegmentDataManager segmentDataManager =
_segmentDataManagerMap.get(segmentName);
if (segmentDataManager == null) {
diff --git
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManager.java
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManager.java
index 934f14a9447..6970d66d105 100644
---
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManager.java
+++
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManager.java
@@ -470,8 +470,7 @@ public class RealtimeTableDataManager extends
BaseTableDataManager {
SegmentZKMetadata zkMetadata = fetchZKMetadata(segmentName);
Preconditions.checkState(zkMetadata.getStatus() != Status.IN_PROGRESS,
"Segment: %s of table: %s is not committed, cannot make it ONLINE",
segmentName, _tableNameWithType);
- IndexLoadingConfig indexLoadingConfig = fetchIndexLoadingConfig();
- indexLoadingConfig.setSegmentTier(zkMetadata.getTier());
+ IndexLoadingConfig indexLoadingConfig =
getCachedIndexLoadingConfig().withSegmentTier(zkMetadata.getTier());
handleSegmentPreload(zkMetadata, indexLoadingConfig);
SegmentDataManager segmentDataManager =
_segmentDataManagerMap.get(segmentName);
if (segmentDataManager == null) {
@@ -541,7 +540,7 @@ public class RealtimeTableDataManager extends
BaseTableDataManager {
_logger.warn("Segment: {} is already completed, skipping adding it as
CONSUMING segment", segmentName);
return;
}
- IndexLoadingConfig indexLoadingConfig = fetchIndexLoadingConfig();
+ IndexLoadingConfig indexLoadingConfig = getCachedIndexLoadingConfig();
handleSegmentPreload(zkMetadata, indexLoadingConfig);
SegmentDataManager segmentDataManager =
_segmentDataManagerMap.get(segmentName);
if (segmentDataManager != null) {
@@ -752,9 +751,7 @@ public class RealtimeTableDataManager extends
BaseTableDataManager {
String segmentName = zkMetadata.getSegmentName();
_logger.info("Downloading and replacing CONSUMING segment: {} with
committed one", segmentName);
File indexDir = downloadSegment(zkMetadata);
- // Get a new index loading config with latest table config and schema to
load the segment
- IndexLoadingConfig indexLoadingConfig = fetchIndexLoadingConfig();
- indexLoadingConfig.setSegmentTier(zkMetadata.getTier());
+ IndexLoadingConfig indexLoadingConfig =
getCachedIndexLoadingConfig().withSegmentTier(zkMetadata.getTier());
addSegment(ImmutableSegmentLoader.load(indexDir, indexLoadingConfig,
_segmentOperationsThrottlerSet, zkMetadata),
zkMetadata);
_ingestionDelayTracker.markPartitionForVerification(segmentName);
@@ -774,8 +771,7 @@ public class RealtimeTableDataManager extends
BaseTableDataManager {
throws Exception {
_logger.info("Replacing CONSUMING segment: {} with the one sealed
locally", segmentName);
File indexDir = new File(_indexDir, segmentName);
- // Get a new index loading config with latest table config and schema to
load the segment
- IndexLoadingConfig indexLoadingConfig = fetchIndexLoadingConfig();
+ IndexLoadingConfig indexLoadingConfig = getCachedIndexLoadingConfig();
ImmutableSegment immutableSegment =
ImmutableSegmentLoader.load(indexDir, indexLoadingConfig,
_segmentOperationsThrottlerSet, zkMetadata);
diff --git
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
index c1fe11501e1..c7faa46cd1a 100644
---
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
+++
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
@@ -30,16 +30,20 @@ import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import org.apache.commons.io.FileUtils;
import org.apache.helix.HelixManager;
+import org.apache.pinot.common.metadata.ZKMetadataProvider;
import org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
import org.apache.pinot.common.metrics.ServerMetrics;
import org.apache.pinot.common.tier.TierFactory;
import org.apache.pinot.common.utils.TarCompressionUtils;
import org.apache.pinot.common.utils.fetcher.BaseSegmentFetcher;
import org.apache.pinot.common.utils.fetcher.SegmentFetcherFactory;
+import org.apache.pinot.common.utils.helix.FakePropertyStore;
import org.apache.pinot.core.data.manager.offline.ImmutableSegmentDataManager;
import org.apache.pinot.core.data.manager.offline.OfflineTableDataManager;
import org.apache.pinot.segment.local.data.manager.SegmentDataManager;
@@ -57,13 +61,17 @@ import org.apache.pinot.segment.spi.SegmentMetadata;
import org.apache.pinot.segment.spi.V1Constants;
import org.apache.pinot.segment.spi.creator.SegmentGeneratorConfig;
import org.apache.pinot.segment.spi.creator.SegmentVersion;
+import org.apache.pinot.segment.spi.index.StandardIndexes;
import org.apache.pinot.segment.spi.index.metadata.SegmentMetadataImpl;
import org.apache.pinot.segment.spi.store.SegmentDirectory;
import org.apache.pinot.segment.spi.store.SegmentDirectoryPaths;
import org.apache.pinot.spi.config.instance.InstanceDataManagerConfig;
+import org.apache.pinot.spi.config.table.FieldConfig;
import org.apache.pinot.spi.config.table.TableConfig;
import org.apache.pinot.spi.config.table.TableType;
import org.apache.pinot.spi.config.table.TierConfig;
+import org.apache.pinot.spi.config.table.TimestampConfig;
+import org.apache.pinot.spi.config.table.TimestampIndexGranularity;
import org.apache.pinot.spi.crypt.PinotCrypter;
import org.apache.pinot.spi.crypt.PinotCrypterFactory;
import org.apache.pinot.spi.data.FieldSpec.DataType;
@@ -88,6 +96,7 @@ import static org.mockito.Mockito.never;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
import static org.testng.Assert.*;
@@ -279,6 +288,28 @@ public class BaseTableDataManagerTest {
assertEquals(new SegmentMetadataImpl(tierDataDir).getTotalDocs(), 5);
}
+ /// Regression test for https://github.com/apache/pinot/issues/18164: the
table-level config passed into reload is
+ /// shared by all segments, so reloading one segment must not set its tier
on the shared config.
+ @Test
+ public void
testReloadSegmentFromDefaultTierDoesNotMutateSharedIndexLoadingConfig()
+ throws Exception {
+ // The current tier is null while the target tier is coolTier, exercising
both tier mutations in the reload path.
+ SegmentZKMetadata zkMetadata = createRawSegment(SegmentVersion.v3, 5);
+ zkMetadata.setTier(TIER_NAME);
+ SegmentMetadata localMetadata = mock(SegmentMetadata.class);
+ when(localMetadata.getCrc()).thenReturn(0L);
+
+ ImmutableSegmentDataManager segmentDataManager =
createImmutableSegmentDataManager(SEGMENT_NAME, localMetadata);
+ BaseTableDataManager tableDataManager = spy(createTableManager());
+ tableDataManager.registerSegment(SEGMENT_NAME, segmentDataManager);
+ seedZKMetadata(tableDataManager, SEGMENT_NAME, zkMetadata);
+
+ IndexLoadingConfig sharedConfig = new
IndexLoadingConfig(DEFAULT_TABLE_CONFIG, SCHEMA);
+ tableDataManager.reloadSegment(segmentDataManager, sharedConfig, true);
+ assertNull(sharedConfig.getSegmentTier());
+ assertNull(sharedConfig.getTableDataDir());
+ }
+
@Test
public void testReloadSegmentConvertVersion()
throws Exception {
@@ -1069,6 +1100,127 @@ public class BaseTableDataManagerTest {
verify(segmentDirectory, times(1)).onSegmentAdded();
}
+ @Test
+ public void testGetCachedIndexLoadingConfigReusesSchemaAndRefreshes() {
+ TableConfig table = new
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME).build();
+ Schema original = createSchemaReuseSchema();
+ BaseTableDataManager manager = createSchemaReuseManager(table, original);
+ IndexLoadingConfig first = manager.getCachedIndexLoadingConfig();
+ IndexLoadingConfig second = manager.getCachedIndexLoadingConfig();
+ assertSame(first.getSchema(), original);
+ assertSame(second.getSchema(), original);
+ assertSame(first.getTableConfig(), second.getTableConfig());
+ assertSame(first, second);
+ verifyNoInteractions(manager._propertyStore);
+
+ Schema changed = createSchemaReuseSchema();
+ changed.getFieldSpecFor("id").setDefaultNullValue(-2);
+ ZKMetadataProvider.setSchema(manager._propertyStore, changed);
+ // Explicit refreshes still fetch the latest schema even without a
separate refresh message.
+ Schema refreshed = manager.fetchIndexLoadingConfig().getSchema();
+ assertNotSame(refreshed, original);
+ assertEquals(refreshed.getFieldSpecFor("id").getDefaultNullValue(), -2);
+ assertEquals(original.getFieldSpecFor("id").getDefaultNullValue(), -1);
+ assertSame(manager.getCachedTableConfigAndSchema().getRight(), refreshed);
+ assertSame(manager.getCachedIndexLoadingConfig().getSchema(), refreshed);
+
+ manager.updateCachedTableConfigAndSchema(table, changed);
+ assertSame(manager.getCachedIndexLoadingConfig().getSchema(), changed);
+ TableConfig indexedTable = new
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME)
+ .setInvertedIndexColumns(List.of("id")).build();
+ manager.updateCachedTableConfigAndSchema(indexedTable, changed);
+ IndexLoadingConfig indexed = manager.getCachedIndexLoadingConfig();
+
assertTrue(indexed.getFieldIndexConfig("id").getConfig(StandardIndexes.inverted()).isEnabled());
+
assertFalse(first.getFieldIndexConfig("id").getConfig(StandardIndexes.inverted()).isEnabled());
+ assertSame(indexed.getFieldIndexConfig("id"),
manager.getCachedIndexLoadingConfig().getFieldIndexConfig("id"));
+ }
+
+ @Test
+ public void testGetCachedIndexLoadingConfigDerivesSegmentTier() {
+ TableConfig table = new
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME).build();
+ Schema schema = createSchemaReuseSchema();
+ BaseTableDataManager manager = createSchemaReuseManager(table, schema);
+ IndexLoadingConfig shared = manager.getCachedIndexLoadingConfig();
+ IndexLoadingConfig tierConfig = shared.withSegmentTier("cold");
+ assertNotSame(tierConfig, shared);
+ assertEquals(tierConfig.getSegmentTier(), "cold");
+ assertNull(shared.getSegmentTier());
+ assertSame(manager.getCachedIndexLoadingConfig(), shared);
+ assertSame(shared.withSegmentTier(null), shared);
+ verifyNoInteractions(manager._propertyStore);
+ }
+
+ @Test
+ public void testGetCachedIndexLoadingConfigNormalizesTimestampBeforeReuse() {
+ BaseTableDataManager manager =
+
createSchemaReuseManager(createTimestampTable(TimestampIndexGranularity.DAY),
createSchemaReuseSchema());
+ Schema first = manager.getCachedIndexLoadingConfig().getSchema();
+ IndexLoadingConfig second = manager.getCachedIndexLoadingConfig();
+ assertSame(second.getSchema(), first);
+ assertTrue(first.hasColumn("$ts$DAY"));
+
assertTrue(second.getFieldIndexConfigByColName().get("$ts$DAY").getConfig(StandardIndexes.range()).isEnabled());
+
assertEquals(second.getTableConfig().getIndexingConfig().getRangeIndexColumns(),
List.of("$ts$DAY"));
+
assertEquals(second.getTableConfig().getIngestionConfig().getTransformConfigs().size(),
1);
+
+ ZKMetadataProvider.setTableConfig(manager._propertyStore,
createTimestampTable(TimestampIndexGranularity.HOUR));
+ manager.onTableConfigOrSchemaRefresh();
+ Schema changed = manager.getCachedIndexLoadingConfig().getSchema();
+ assertNotSame(changed, first);
+ assertTrue(changed.hasColumn("$ts$HOUR"));
+ assertFalse(changed.hasColumn("$ts$DAY"));
+ assertFalse(first.hasColumn("$ts$HOUR"));
+ }
+
+ @Test
+ public void testGetCachedIndexLoadingConfigConcurrentlyReusesCachedSchema()
+ throws Exception {
+ BaseTableDataManager manager =
+
createSchemaReuseManager(createTimestampTable(TimestampIndexGranularity.DAY),
createSchemaReuseSchema());
+ IndexLoadingConfig shared = manager.getCachedIndexLoadingConfig();
+ ExecutorService executor = Executors.newFixedThreadPool(8);
+ CountDownLatch start = new CountDownLatch(1);
+ try {
+ List<Future<IndexLoadingConfig>> results = new ArrayList<>();
+ for (int i = 0; i < 32; i++) {
+ results.add(executor.submit(() -> {
+ assertTrue(start.await(10, TimeUnit.SECONDS));
+ return manager.getCachedIndexLoadingConfig();
+ }));
+ }
+ start.countDown();
+ for (Future<IndexLoadingConfig> result : results) {
+ assertSame(result.get(10, TimeUnit.SECONDS), shared);
+ }
+ verifyNoInteractions(manager._propertyStore);
+ } finally {
+ start.countDown();
+ executor.shutdownNow();
+ }
+ }
+
+ private static BaseTableDataManager createSchemaReuseManager(TableConfig
table, Schema schema) {
+ BaseTableDataManager manager = new OfflineTableDataManager();
+ manager._propertyStore = new FakePropertyStore();
+ manager._tableNameWithType = OFFLINE_TABLE_NAME;
+ ZKMetadataProvider.setTableConfig(manager._propertyStore, table);
+ ZKMetadataProvider.setSchema(manager._propertyStore, schema);
+ manager.updateCachedTableConfigAndSchema(table, schema);
+ manager._propertyStore = spy(manager._propertyStore);
+ return manager;
+ }
+
+ private static Schema createSchemaReuseSchema() {
+ return new Schema.SchemaBuilder().setSchemaName(RAW_TABLE_NAME)
+ .addSingleValueDimension("id", DataType.INT, -1)
+ .addDateTime("ts", DataType.TIMESTAMP, "TIMESTAMP",
"1:MILLISECONDS").build();
+ }
+
+ private static TableConfig createTimestampTable(TimestampIndexGranularity
granularity) {
+ return new
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME)
+ .setFieldConfigList(List.of(new FieldConfig.Builder("ts")
+ .withTimestampConfig(new
TimestampConfig(List.of(granularity))).build())).build();
+ }
+
protected BaseTableDataManager createTableManager() {
return createTableManager(createDefaultInstanceDataManagerConfig());
}
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentLoader.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentLoader.java
index 8f3f4e62dc7..c593f69a58f 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentLoader.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentLoader.java
@@ -127,6 +127,7 @@ public class ImmutableSegmentLoader {
if (segmentMetadata.getTotalDocs() == 0) {
return new EmptyIndexSegment(segmentMetadata);
}
+ indexLoadingConfig =
indexLoadingConfig.withOpenStructChildConfigs(segmentMetadata);
String segmentName = segmentMetadata.getName();
SegmentDirectoryLoaderContext segmentLoaderContext = new
SegmentDirectoryLoaderContext.Builder()
.setReadMode(indexLoadingConfig.getReadMode())
@@ -171,6 +172,7 @@ public class ImmutableSegmentLoader {
SegmentMetadataImpl segmentMetadata = new SegmentMetadataImpl(indexDir);
if (segmentMetadata.getTotalDocs() > 0) {
+ indexLoadingConfig =
indexLoadingConfig.withOpenStructChildConfigs(segmentMetadata);
if (segmentOperationsThrottlerSet != null) {
segmentOperationsThrottlerSet.getSegmentAllIndexPreprocessThrottler().acquire();
}
@@ -205,6 +207,7 @@ public class ImmutableSegmentLoader {
// mirroring the non-empty ImmutableSegmentImpl path.
return new EmptyIndexSegment(segmentMetadata, segmentDirectory);
}
+ indexLoadingConfig =
indexLoadingConfig.withOpenStructChildConfigs(segmentMetadata);
// Remove columns not in schema from the metadata
Map<String, ColumnMetadata> columnMetadataMap =
segmentMetadata.getColumnMetadataMap();
@@ -229,7 +232,7 @@ public class ImmutableSegmentLoader {
}
}
} else {
- indexLoadingConfig.addKnownColumns(columnMetadataMap.keySet());
+ indexLoadingConfig =
indexLoadingConfig.withKnownColumns(columnMetadataMap.keySet());
}
SegmentDirectory.Reader segmentReader = segmentDirectory.createReader();
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/loader/IndexLoadingConfig.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/loader/IndexLoadingConfig.java
index 3782f3abbf8..e68480ace43 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/loader/IndexLoadingConfig.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/loader/IndexLoadingConfig.java
@@ -24,6 +24,7 @@ import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
import java.util.Set;
import javax.annotation.Nullable;
import org.apache.commons.lang3.StringUtils;
@@ -44,64 +45,200 @@ import
org.apache.pinot.spi.config.table.MultiColumnTextIndexConfig;
import org.apache.pinot.spi.config.table.OpenStructIndexConfig;
import org.apache.pinot.spi.config.table.StarTreeIndexConfig;
import org.apache.pinot.spi.config.table.TableConfig;
+import org.apache.pinot.spi.data.ComplexFieldSpec;
import org.apache.pinot.spi.data.DimensionFieldSpec;
import org.apache.pinot.spi.data.FieldSpec;
+import org.apache.pinot.spi.data.FieldSpec.DataType;
import org.apache.pinot.spi.data.OpenStructNaming;
import org.apache.pinot.spi.data.Schema;
import org.apache.pinot.spi.utils.ReadMode;
import org.apache.pinot.spi.utils.TimestampIndexUtils;
-/// Table level index loading config.
+/// Index loading config with shared table-level state and segment-local
mutable overrides.
public class IndexLoadingConfig {
private static final int DEFAULT_REALTIME_AVG_MULTI_VALUE_COUNT = 2;
public static final String READ_MODE_KEY = "readMode";
- private final InstanceDataManagerConfig _instanceDataManagerConfig;
- private final TableConfig _tableConfig;
- private final Schema _schema;
+ private final ImmutableState _immutableState;
- // These fields can be modified after initialization
- // TODO: Revisit them
- private ReadMode _readMode = ReadMode.DEFAULT_MODE;
- private SegmentVersion _segmentVersion;
+ // Mutable config and segment-specific overrides.
+ @Nullable
+ private ReadMode _readModeOverride;
+ @Nullable
+ private SegmentVersion _segmentVersionOverride;
private String _segmentTier;
private Set<String> _knownColumns;
private String _tableDataDir;
private boolean _errorOnColumnBuildFailure;
private boolean _forwardIndexOnly;
+ private ResolvedIndexState _resolvedIndexState;
+
+ /// Immutable table-level state shared by derived segment configs.
+ private static final class ImmutableState {
+ @Nullable
+ private final InstanceDataManagerConfig _instanceDataManagerConfig;
+ @Nullable
+ private final TableConfig _tableConfig;
+ @Nullable
+ private final Schema _schema;
+ private final ReadMode _readMode;
+ @Nullable
+ private final SegmentVersion _segmentVersion;
+ @Nullable
+ private final String _instanceId;
+ private final boolean _isRealtimeOffHeapAllocation;
+ private final boolean _isDirectRealtimeOffHeapAllocation;
+ private final int _realtimeAvgMultiValueCount;
+ @Nullable
+ private final String _segmentStoreURI;
+ @Nullable
+ private final String _segmentDirectoryLoader;
+ @Nullable
+ private final Map<String, Map<String, String>> _instanceTierConfigs;
+ private final List<String> _sortedColumns;
+ private final ColumnMinMaxValueGeneratorMode
_columnMinMaxValueGeneratorMode;
+ private final boolean _hasOpenStructColumns;
+
+ private ImmutableState(@Nullable InstanceDataManagerConfig
instanceDataManagerConfig,
+ @Nullable TableConfig tableConfig, @Nullable Schema schema) {
+ _instanceDataManagerConfig = instanceDataManagerConfig;
+ _tableConfig = tableConfig;
+ _schema = schema;
+
+ String instanceId = null;
+ boolean isRealtimeOffHeapAllocation = false;
+ boolean isDirectRealtimeOffHeapAllocation = false;
+ int realtimeAvgMultiValueCount = DEFAULT_REALTIME_AVG_MULTI_VALUE_COUNT;
+ ReadMode readMode = ReadMode.DEFAULT_MODE;
+ SegmentVersion segmentVersion = null;
+ String segmentStoreURI = null;
+ String segmentDirectoryLoader = null;
+ Map<String, Map<String, String>> instanceTierConfigs = null;
+ if (instanceDataManagerConfig != null) {
+ ReadMode instanceReadMode = instanceDataManagerConfig.getReadMode();
+ if (instanceReadMode != null) {
+ readMode = instanceReadMode;
+ }
+ String instanceSegmentVersion =
instanceDataManagerConfig.getSegmentFormatVersion();
+ if (instanceSegmentVersion != null) {
+ segmentVersion =
SegmentVersion.valueOf(instanceSegmentVersion.toLowerCase());
+ }
+ instanceId = instanceDataManagerConfig.getInstanceId();
+ isRealtimeOffHeapAllocation =
instanceDataManagerConfig.isRealtimeOffHeapAllocation();
+ isDirectRealtimeOffHeapAllocation =
instanceDataManagerConfig.isDirectRealtimeOffHeapAllocation();
+ String avgMultiValueCount =
instanceDataManagerConfig.getAvgMultiValueCount();
+ if (avgMultiValueCount != null) {
+ realtimeAvgMultiValueCount = Integer.parseInt(avgMultiValueCount);
+ }
+ segmentStoreURI = instanceDataManagerConfig.getSegmentStoreUri();
+ segmentDirectoryLoader =
instanceDataManagerConfig.getSegmentDirectoryLoader();
+ Map<String, Map<String, String>> tierConfigs =
instanceDataManagerConfig.getTierConfigs();
+ instanceTierConfigs = tierConfigs != null ? tierConfigs : Map.of();
+ }
- // Initialized by instance data manager config
- private String _instanceId;
- private boolean _isRealtimeOffHeapAllocation;
- private boolean _isDirectRealtimeOffHeapAllocation;
- private int _realtimeAvgMultiValueCount =
DEFAULT_REALTIME_AVG_MULTI_VALUE_COUNT;
- private String _segmentStoreURI;
- private String _segmentDirectoryLoader;
- private Map<String, Map<String, String>> _instanceTierConfigs;
+ List<String> sortedColumns = List.of();
+ ColumnMinMaxValueGeneratorMode columnMinMaxValueGeneratorMode =
ColumnMinMaxValueGeneratorMode.DEFAULT_MODE;
+ boolean hasOpenStructColumns = false;
+ if (tableConfig != null) {
+ if (schema != null) {
+ TimestampIndexUtils.applyTimestampIndex(tableConfig, schema);
+ for (ComplexFieldSpec fieldSpec : schema.getComplexFieldSpecs()) {
+ if (fieldSpec.getDataType() == DataType.OPEN_STRUCT) {
+ hasOpenStructColumns = true;
+ break;
+ }
+ }
+ }
+ IndexingConfig indexingConfig = tableConfig.getIndexingConfig();
+ String tableReadMode = indexingConfig.getLoadMode();
+ if (tableReadMode != null) {
+ readMode = ReadMode.getEnum(tableReadMode);
+ }
+ String tableSegmentVersion = indexingConfig.getSegmentFormatVersion();
+ if (tableSegmentVersion != null) {
+ segmentVersion =
SegmentVersion.valueOf(tableSegmentVersion.toLowerCase());
+ }
+ List<String> tableSortedColumns = indexingConfig.getSortedColumn();
+ if (tableSortedColumns != null) {
+ sortedColumns = tableSortedColumns;
+ }
+ String generatorMode =
indexingConfig.getColumnMinMaxValueGeneratorMode();
+ if (generatorMode != null) {
+ columnMinMaxValueGeneratorMode =
ColumnMinMaxValueGeneratorMode.valueOf(generatorMode.toUpperCase());
+ }
+ }
- // Initialized by table config and schema
- private List<String> _sortedColumns = List.of();
- private ColumnMinMaxValueGeneratorMode _columnMinMaxValueGeneratorMode =
ColumnMinMaxValueGeneratorMode.DEFAULT_MODE;
- private boolean _enableDynamicStarTreeCreation;
- private List<StarTreeIndexConfig> _starTreeIndexConfigs;
- private boolean _enableDefaultStarTree;
- private Map<String, FieldIndexConfigs> _indexConfigsByColName = new
HashMap<>();
- private boolean _skipSegmentPreprocess;
+ _instanceId = instanceId;
+ _readMode = readMode;
+ _segmentVersion = segmentVersion;
+ _isRealtimeOffHeapAllocation = isRealtimeOffHeapAllocation;
+ _isDirectRealtimeOffHeapAllocation = isDirectRealtimeOffHeapAllocation;
+ _realtimeAvgMultiValueCount = realtimeAvgMultiValueCount;
+ _segmentStoreURI = segmentStoreURI;
+ _segmentDirectoryLoader = segmentDirectoryLoader;
+ _instanceTierConfigs = instanceTierConfigs;
+ _sortedColumns = sortedColumns;
+ _columnMinMaxValueGeneratorMode = columnMinMaxValueGeneratorMode;
+ _hasOpenStructColumns = hasOpenStructColumns;
+ }
+ }
- private boolean _dirty = true;
+ /// Index settings resolved from the table config, segment tier, schema, and
known segment columns.
+ private static final class ResolvedIndexState {
+ private static final ResolvedIndexState EMPTY =
+ new ResolvedIndexState(false, null, false, Map.of(), false, null);
+
+ private final boolean _enableDynamicStarTreeCreation;
+ @Nullable
+ private final List<StarTreeIndexConfig> _starTreeIndexConfigs;
+ private final boolean _enableDefaultStarTree;
+ private final Map<String, FieldIndexConfigs> _indexConfigsByColName;
+ private final boolean _skipSegmentPreprocess;
+ @Nullable
+ private final MultiColumnTextIndexConfig _multiColTextIndexConfig;
+
+ private ResolvedIndexState(boolean enableDynamicStarTreeCreation,
+ @Nullable List<StarTreeIndexConfig> starTreeIndexConfigs, boolean
enableDefaultStarTree,
+ Map<String, FieldIndexConfigs> indexConfigsByColName, boolean
skipSegmentPreprocess,
+ @Nullable MultiColumnTextIndexConfig multiColTextIndexConfig) {
+ _enableDynamicStarTreeCreation = enableDynamicStarTreeCreation;
+ _starTreeIndexConfigs = starTreeIndexConfigs;
+ _enableDefaultStarTree = enableDefaultStarTree;
+ _indexConfigsByColName = indexConfigsByColName;
+ _skipSegmentPreprocess = skipSegmentPreprocess;
+ _multiColTextIndexConfig = multiColTextIndexConfig;
+ }
- private MultiColumnTextIndexConfig _multiColTextIndexConfig;
+ private ResolvedIndexState withIndexConfigsByColName(Map<String,
FieldIndexConfigs> indexConfigsByColName) {
+ return new ResolvedIndexState(_enableDynamicStarTreeCreation,
_starTreeIndexConfigs, _enableDefaultStarTree,
+ indexConfigsByColName, _skipSegmentPreprocess,
_multiColTextIndexConfig);
+ }
+ }
/// NOTE: This step might modify the passed in table config and schema.
///
/// TODO: Revisit the init handling. Currently it doesn't apply tiered
config override
public IndexLoadingConfig(@Nullable InstanceDataManagerConfig
instanceDataManagerConfig,
@Nullable TableConfig tableConfig, @Nullable Schema schema) {
- _instanceDataManagerConfig = instanceDataManagerConfig;
- _tableConfig = tableConfig;
- _schema = schema;
- init();
+ _immutableState = new ImmutableState(instanceDataManagerConfig,
tableConfig, schema);
+ if (tableConfig != null) {
+ refreshIndexConfigs();
+ }
+ }
+
+ /// Creates a segment-local wrapper around the already processed table-level
config. The wrapper shares the
+ /// table config, schema, and resolved index configs until a
segment-specific override requires a local copy.
+ private IndexLoadingConfig(IndexLoadingConfig source) {
+ _immutableState = source._immutableState;
+ _readModeOverride = source._readModeOverride;
+ _segmentVersionOverride = source._segmentVersionOverride;
+ _segmentTier = source._segmentTier;
+ _knownColumns = source._knownColumns;
+ _tableDataDir = source._tableDataDir;
+ _errorOnColumnBuildFailure = source._errorOnColumnBuildFailure;
+ _forwardIndexOnly = source._forwardIndexOnly;
+ _resolvedIndexState = source._resolvedIndexState;
}
@VisibleForTesting
@@ -121,88 +258,22 @@ public class IndexLoadingConfig {
@Nullable
public InstanceDataManagerConfig getInstanceDataManagerConfig() {
- return _instanceDataManagerConfig;
+ return _immutableState._instanceDataManagerConfig;
}
@Nullable
public TableConfig getTableConfig() {
- return _tableConfig;
+ return _immutableState._tableConfig;
}
@Nullable
public Schema getSchema() {
- return _schema;
- }
-
- private void init() {
- if (_instanceDataManagerConfig != null) {
- extractFromInstanceConfig();
- }
- if (_tableConfig != null) {
- extractFromTableConfigAndSchema();
- }
- }
-
- private void extractFromInstanceConfig() {
- _instanceId = _instanceDataManagerConfig.getInstanceId();
-
- ReadMode instanceReadMode = _instanceDataManagerConfig.getReadMode();
- if (instanceReadMode != null) {
- _readMode = instanceReadMode;
- }
-
- String instanceSegmentVersion =
_instanceDataManagerConfig.getSegmentFormatVersion();
- if (instanceSegmentVersion != null) {
- _segmentVersion =
SegmentVersion.valueOf(instanceSegmentVersion.toLowerCase());
- }
-
- _isRealtimeOffHeapAllocation =
_instanceDataManagerConfig.isRealtimeOffHeapAllocation();
- _isDirectRealtimeOffHeapAllocation =
_instanceDataManagerConfig.isDirectRealtimeOffHeapAllocation();
-
- String avgMultiValueCount =
_instanceDataManagerConfig.getAvgMultiValueCount();
- if (avgMultiValueCount != null) {
- _realtimeAvgMultiValueCount = Integer.parseInt(avgMultiValueCount);
- }
- _segmentStoreURI = _instanceDataManagerConfig.getSegmentStoreUri();
- _segmentDirectoryLoader =
_instanceDataManagerConfig.getSegmentDirectoryLoader();
-
- Map<String, Map<String, String>> tierConfigs =
_instanceDataManagerConfig.getTierConfigs();
- _instanceTierConfigs = tierConfigs != null ? tierConfigs : Map.of();
- }
-
- private void extractFromTableConfigAndSchema() {
- if (_schema != null) {
- TimestampIndexUtils.applyTimestampIndex(_tableConfig, _schema);
- }
-
- IndexingConfig indexingConfig = _tableConfig.getIndexingConfig();
- String tableReadMode = indexingConfig.getLoadMode();
- if (tableReadMode != null) {
- _readMode = ReadMode.getEnum(tableReadMode);
- }
-
- List<String> sortedColumns = indexingConfig.getSortedColumn();
- if (sortedColumns != null) {
- _sortedColumns = sortedColumns;
- }
-
- String tableSegmentVersion = indexingConfig.getSegmentFormatVersion();
- if (tableSegmentVersion != null) {
- _segmentVersion =
SegmentVersion.valueOf(tableSegmentVersion.toLowerCase());
- }
-
- String columnMinMaxValueGeneratorMode =
indexingConfig.getColumnMinMaxValueGeneratorMode();
- if (columnMinMaxValueGeneratorMode != null) {
- _columnMinMaxValueGeneratorMode =
-
ColumnMinMaxValueGeneratorMode.valueOf(columnMinMaxValueGeneratorMode.toUpperCase());
- }
-
- refreshIndexConfigs();
+ return _immutableState._schema;
}
public void refreshIndexConfigs() {
- if (_tableConfig == null) {
- _dirty = false;
+ if (_immutableState._tableConfig == null) {
+ _resolvedIndexState = ResolvedIndexState.EMPTY;
return;
}
// Accessing the index configs for single-column index is handled by
IndexType.getConfig() as defined in index-spi.
@@ -210,112 +281,105 @@ public class IndexLoadingConfig {
// specific index configs transparently.
TableConfig tableConfig = getTableConfigWithTierOverwrites();
Schema schema = inferSchema();
- _indexConfigsByColName =
FieldIndexConfigsUtil.createIndexConfigsByColName(tableConfig, schema);
+ Map<String, FieldIndexConfigs> indexConfigsByColName =
+ FieldIndexConfigsUtil.createIndexConfigsByColName(tableConfig, schema);
// Accessing the StarTree index configs is not handled by
IndexType.getConfig(), so we manually update them.
IndexingConfig indexingConfig = tableConfig.getIndexingConfig();
- _enableDynamicStarTreeCreation =
indexingConfig.isEnableDynamicStarTreeCreation();
- _starTreeIndexConfigs = indexingConfig.getStarTreeIndexConfigs();
- _enableDefaultStarTree = indexingConfig.isEnableDefaultStarTree();
- _multiColTextIndexConfig = indexingConfig.getMultiColumnTextIndexConfig();
- _skipSegmentPreprocess = indexingConfig.isSkipSegmentPreprocess();
- _dirty = false;
+ _resolvedIndexState = new
ResolvedIndexState(indexingConfig.isEnableDynamicStarTreeCreation(),
+ indexingConfig.getStarTreeIndexConfigs(),
indexingConfig.isEnableDefaultStarTree(), indexConfigsByColName,
+ indexingConfig.isSkipSegmentPreprocess(),
indexingConfig.getMultiColumnTextIndexConfig());
+ }
+
+ private ResolvedIndexState getResolvedIndexState() {
+ if (_resolvedIndexState == null) {
+ refreshIndexConfigs();
+ }
+ return _resolvedIndexState;
}
private TableConfig getTableConfigWithTierOverwrites() {
- return (_segmentTier == null || _tableConfig == null) ? _tableConfig
- : TableConfigUtils.overwriteTableConfigForTier(_tableConfig,
_segmentTier);
+ return _segmentTier == null || _immutableState._tableConfig == null ?
_immutableState._tableConfig
+ :
TableConfigUtils.overwriteTableConfigForTier(_immutableState._tableConfig,
_segmentTier);
}
private Schema inferSchema() {
- if (_schema != null) {
- return _schema;
+ if (_immutableState._schema != null) {
+ return _immutableState._schema;
}
Schema schema = new Schema();
for (String column : getAllKnownColumns()) {
- schema.addField(new DimensionFieldSpec(column,
FieldSpec.DataType.STRING, true));
+ schema.addField(new DimensionFieldSpec(column, DataType.STRING, true));
}
return schema;
}
public ReadMode getReadMode() {
- return _readMode;
+ return _readModeOverride != null ? _readModeOverride :
_immutableState._readMode;
}
public void setReadMode(ReadMode readMode) {
- _readMode = readMode;
+ _readModeOverride = readMode;
}
public List<String> getSortedColumns() {
- return unmodifiable(_sortedColumns);
+ return unmodifiable(_immutableState._sortedColumns);
}
public boolean isEnableDynamicStarTreeCreation() {
- if (_dirty) {
- refreshIndexConfigs();
- }
- return _enableDynamicStarTreeCreation;
+ return getResolvedIndexState()._enableDynamicStarTreeCreation;
}
@Nullable
public List<StarTreeIndexConfig> getStarTreeIndexConfigs() {
- if (_dirty) {
- refreshIndexConfigs();
- }
- return unmodifiable(_starTreeIndexConfigs);
+ return unmodifiable(getResolvedIndexState()._starTreeIndexConfigs);
}
@Nullable
public MultiColumnTextIndexConfig getMultiColTextIndexConfig() {
- if (_dirty) {
- refreshIndexConfigs();
- }
- return _multiColTextIndexConfig;
+ return getResolvedIndexState()._multiColTextIndexConfig;
}
public boolean isEnableDefaultStarTree() {
- if (_dirty) {
- refreshIndexConfigs();
- }
- return _enableDefaultStarTree;
+ return getResolvedIndexState()._enableDefaultStarTree;
}
@Nullable
public SegmentVersion getSegmentVersion() {
- return _segmentVersion;
+ return _segmentVersionOverride != null ? _segmentVersionOverride :
_immutableState._segmentVersion;
}
/// For tests only.
public void setSegmentVersion(SegmentVersion segmentVersion) {
- _segmentVersion = segmentVersion;
+ _segmentVersionOverride = segmentVersion;
}
public boolean isRealtimeOffHeapAllocation() {
- return _isRealtimeOffHeapAllocation;
+ return _immutableState._isRealtimeOffHeapAllocation;
}
public boolean isDirectRealtimeOffHeapAllocation() {
- return _isDirectRealtimeOffHeapAllocation;
+ return _immutableState._isDirectRealtimeOffHeapAllocation;
}
public ColumnMinMaxValueGeneratorMode getColumnMinMaxValueGeneratorMode() {
- return _columnMinMaxValueGeneratorMode;
+ return _immutableState._columnMinMaxValueGeneratorMode;
}
public String getSegmentStoreURI() {
- return _segmentStoreURI;
+ return _immutableState._segmentStoreURI;
}
public int getRealtimeAvgMultiValueCount() {
- return _realtimeAvgMultiValueCount;
+ return _immutableState._realtimeAvgMultiValueCount;
}
public String getSegmentDirectoryLoader() {
- return StringUtils.isNotBlank(_segmentDirectoryLoader) ?
_segmentDirectoryLoader
+ return StringUtils.isNotBlank(_immutableState._segmentDirectoryLoader) ?
_immutableState._segmentDirectoryLoader
: SegmentDirectoryLoaderRegistry.DEFAULT_SEGMENT_DIRECTORY_LOADER_NAME;
}
public String getInstanceId() {
- return _instanceId;
+ return _immutableState._instanceId;
}
public String getSegmentTier() {
@@ -323,8 +387,26 @@ public class IndexLoadingConfig {
}
public void setSegmentTier(String segmentTier) {
+ if (Objects.equals(_segmentTier, segmentTier)) {
+ return;
+ }
_segmentTier = segmentTier;
- _dirty = true;
+ _resolvedIndexState = null;
+ }
+
+ public IndexLoadingConfig withSegmentTier(@Nullable String segmentTier) {
+ if (Objects.equals(_segmentTier, segmentTier)) {
+ return this;
+ }
+ return copyWithSegmentTier(segmentTier);
+ }
+
+ /// Creates a segment-local mutable wrapper with the given tier. This always
creates a wrapper, even when the tier is
+ /// unchanged, so callers can safely apply additional segment-specific
overrides.
+ public IndexLoadingConfig copyWithSegmentTier(@Nullable String segmentTier) {
+ IndexLoadingConfig derived = new IndexLoadingConfig(this);
+ derived.setSegmentTier(segmentTier);
+ return derived;
}
public String getTableDataDir() {
@@ -352,25 +434,16 @@ public class IndexLoadingConfig {
}
public boolean isSkipSegmentPreprocess() {
- if (_dirty) {
- refreshIndexConfigs();
- }
- return _skipSegmentPreprocess;
+ return getResolvedIndexState()._skipSegmentPreprocess;
}
@Nullable
public FieldIndexConfigs getFieldIndexConfig(String columnName) {
- if (_indexConfigsByColName == null || _dirty) {
- refreshIndexConfigs();
- }
- return _indexConfigsByColName.get(columnName);
+ return getResolvedIndexState()._indexConfigsByColName.get(columnName);
}
public Map<String, FieldIndexConfigs> getFieldIndexConfigByColName() {
- if (_indexConfigsByColName == null || _dirty) {
- refreshIndexConfigs();
- }
- return unmodifiable(_indexConfigsByColName);
+ return unmodifiable(getResolvedIndexState()._indexConfigsByColName);
}
/// Returns a subset of the columns on the table.
@@ -379,10 +452,11 @@ public class IndexLoadingConfig {
/// tries its bests to get the columns from other attributes like
[#getTableConfig()], which may also not be
/// defined or may not be complete.
private Set<String> getAllKnownColumns() {
- assert _tableConfig != null && _schema == null;
+ assert _immutableState._tableConfig != null && _immutableState._schema ==
null;
if (_knownColumns == null) {
- Set<String> knownColumns =
_tableConfig.getIndexingConfig().getAllReferencedColumns();
- List<FieldConfig> fieldConfigs = _tableConfig.getFieldConfigList();
+ Set<String> knownColumns =
+ new
HashSet<>(_immutableState._tableConfig.getIndexingConfig().getAllReferencedColumns());
+ List<FieldConfig> fieldConfigs =
_immutableState._tableConfig.getFieldConfigList();
if (fieldConfigs != null) {
for (FieldConfig fieldConfig : fieldConfigs) {
knownColumns.add(fieldConfig.getName());
@@ -394,7 +468,7 @@ public class IndexLoadingConfig {
}
public Map<String, Map<String, String>> getInstanceTierConfigs() {
- return unmodifiable(_instanceTierConfigs);
+ return unmodifiable(_immutableState._instanceTierConfigs);
}
private <E> List<E> unmodifiable(List<E> list) {
@@ -409,20 +483,23 @@ public class IndexLoadingConfig {
return map == null ? null : Collections.unmodifiableMap(map);
}
- public void addOpenStructChildConfigs(SegmentMetadataImpl segmentMetadata) {
- if (_indexConfigsByColName == null || _dirty) {
- refreshIndexConfigs();
+ public IndexLoadingConfig withOpenStructChildConfigs(SegmentMetadataImpl
segmentMetadata) {
+ if (!_immutableState._hasOpenStructColumns) {
+ return this;
}
+ ResolvedIndexState resolvedIndexState = getResolvedIndexState();
+ Map<String, FieldIndexConfigs> indexConfigsByColName =
resolvedIndexState._indexConfigsByColName;
+ Map<String, FieldIndexConfigs> updatedConfigs = null;
for (Map.Entry<String, ColumnMetadata> entry :
segmentMetadata.getColumnMetadataMap().entrySet()) {
String childColumn = entry.getKey();
- if (!childColumn.contains(OpenStructNaming.SEPARATOR) ||
_indexConfigsByColName.containsKey(childColumn)) {
+ if (!childColumn.contains(OpenStructNaming.SEPARATOR) ||
indexConfigsByColName.containsKey(childColumn)) {
continue;
}
if (OpenStructNaming.isSparseColumn(childColumn)) {
continue;
}
String parentColumn = OpenStructNaming.parseParentColumn(childColumn);
- FieldIndexConfigs parentConfigs =
_indexConfigsByColName.get(parentColumn);
+ FieldIndexConfigs parentConfigs =
indexConfigsByColName.get(parentColumn);
if (parentConfigs == null) {
continue;
}
@@ -442,16 +519,39 @@ public class IndexLoadingConfig {
FieldIndexConfigsUtil.fromFieldConfig(keyFieldConfig,
childFieldSpec))
.add(StandardIndexes.inverted(), enableInverted ?
IndexConfig.ENABLED : IndexConfig.DISABLED)
.build();
- _indexConfigsByColName.put(childColumn, childConfigs);
+ if (updatedConfigs == null) {
+ updatedConfigs = new HashMap<>(indexConfigsByColName);
+ }
+ updatedConfigs.put(childColumn, childConfigs);
+ }
+ if (updatedConfigs == null) {
+ return this;
+ }
+ IndexLoadingConfig derived = new IndexLoadingConfig(this);
+ derived._resolvedIndexState =
resolvedIndexState.withIndexConfigsByColName(updatedConfigs);
+ return derived;
+ }
+
+ public void addOpenStructChildConfigs(SegmentMetadataImpl segmentMetadata) {
+ _resolvedIndexState =
withOpenStructChildConfigs(segmentMetadata)._resolvedIndexState;
+ }
+
+ public IndexLoadingConfig withKnownColumns(Set<String> columns) {
+ if (_knownColumns != null && _knownColumns.containsAll(columns)) {
+ return this;
}
+ IndexLoadingConfig derived = new IndexLoadingConfig(this);
+ derived.addKnownColumns(columns);
+ return derived;
}
public void addKnownColumns(Set<String> columns) {
- if (_knownColumns == null) {
- _knownColumns = new HashSet<>(columns);
- } else {
- _knownColumns.addAll(columns);
+ if (_knownColumns != null && _knownColumns.containsAll(columns)) {
+ return;
}
- _dirty = true;
+ Set<String> knownColumns = _knownColumns != null ? new
HashSet<>(_knownColumns) : new HashSet<>();
+ knownColumns.addAll(columns);
+ _knownColumns = knownColumns;
+ _resolvedIndexState = null;
}
}
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/loader/IndexLoadingConfigTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/loader/IndexLoadingConfigTest.java
index 4edb1c383cf..4b246e98afc 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/loader/IndexLoadingConfigTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/loader/IndexLoadingConfigTest.java
@@ -23,18 +23,26 @@ import java.io.IOException;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
+import java.util.Set;
+import java.util.TreeMap;
+import org.apache.pinot.segment.spi.ColumnMetadata;
import org.apache.pinot.segment.spi.index.FieldIndexConfigs;
import org.apache.pinot.segment.spi.index.ForwardIndexConfig;
import org.apache.pinot.segment.spi.index.StandardIndexes;
+import org.apache.pinot.segment.spi.index.metadata.SegmentMetadataImpl;
import org.apache.pinot.spi.config.instance.InstanceDataManagerConfig;
import org.apache.pinot.spi.config.table.FieldConfig;
import org.apache.pinot.spi.config.table.StarTreeIndexConfig;
import org.apache.pinot.spi.config.table.TableConfig;
import org.apache.pinot.spi.config.table.TableType;
+import org.apache.pinot.spi.data.ComplexFieldSpec;
+import org.apache.pinot.spi.data.DimensionFieldSpec;
import org.apache.pinot.spi.data.FieldSpec;
+import org.apache.pinot.spi.data.FieldSpec.DataType;
import org.apache.pinot.spi.data.Schema;
import org.apache.pinot.spi.env.PinotConfiguration;
import org.apache.pinot.spi.utils.JsonUtils;
+import org.apache.pinot.spi.utils.ReadMode;
import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
import org.testng.annotations.Test;
@@ -46,6 +54,24 @@ import static org.testng.Assert.*;
public class IndexLoadingConfigTest {
private static final String TABLE_NAME = "table01";
+ @Test
+ public void testReadModePrecedenceAndOverrideIsolation() {
+ InstanceDataManagerConfig instanceConfig =
mock(InstanceDataManagerConfig.class);
+ when(instanceConfig.getReadMode()).thenReturn(ReadMode.heap);
+ TableConfig tableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(TABLE_NAME)
+ .setLoadMode("MMAP").build();
+
+ IndexLoadingConfig base = new IndexLoadingConfig(instanceConfig,
tableConfig, null);
+ assertEquals(base.getReadMode(), ReadMode.mmap);
+
+ IndexLoadingConfig derived = base.withSegmentTier("coldTier");
+ derived.setReadMode(ReadMode.heap);
+ assertEquals(derived.getReadMode(), ReadMode.heap);
+ assertEquals(base.getReadMode(), ReadMode.mmap);
+
+ assertEquals(new IndexLoadingConfig(instanceConfig, null,
null).getReadMode(), ReadMode.heap);
+ }
+
@Test
public void testCalculateIndexConfigsWithoutTierOverwrites()
throws IOException {
@@ -141,8 +167,10 @@ public class IndexLoadingConfigTest {
.setStarTreeIndexConfigs(List.of(stIdxCfg))
.setTierOverwrites(JsonUtils.stringToJsonNode("{\"coldTier\":
{\"starTreeIndexConfigs\": []}}"))
.setFieldConfigList(Arrays.asList(col1Cfg, col2Cfg)).build();
- IndexLoadingConfig ilc = new IndexLoadingConfig(idmCfg, tableConfig,
schema);
- ilc.setSegmentTier("coldTier");
+ IndexLoadingConfig base = new IndexLoadingConfig(idmCfg, tableConfig,
schema);
+ IndexLoadingConfig ilc = base.withSegmentTier("coldTier");
+ assertSame(base.withSegmentTier(null), base);
+ assertNull(base.getSegmentTier());
// Check index configs for coldTier
assertEquals(ilc.getStarTreeIndexConfigs().size(), 0);
Map<String, FieldIndexConfigs> allFieldCfgs =
ilc.getFieldIndexConfigByColName();
@@ -154,6 +182,36 @@ public class IndexLoadingConfigTest {
assertFalse(fieldCfgs.getConfig(StandardIndexes.inverted()).isEnabled());
assertFalse(fieldCfgs.getConfig(StandardIndexes.bloomFilter()).isEnabled());
assertFalse(fieldCfgs.getConfig(StandardIndexes.dictionary()).isEnabled());
+ assertEquals(base.getStarTreeIndexConfigs().size(), 1);
+
assertTrue(base.getFieldIndexConfig("col1").getConfig(StandardIndexes.inverted()).isEnabled());
+
assertTrue(base.getFieldIndexConfig("col2").getConfig(StandardIndexes.dictionary()).isEnabled());
+ }
+
+ @Test
+ public void testDerivedConfigIsolatesSegmentColumns()
+ throws IOException {
+ Schema schema = new Schema.SchemaBuilder().setSchemaName(TABLE_NAME)
+ .addField(new ComplexFieldSpec("event", DataType.OPEN_STRUCT, true,
Map.of())).build();
+ FieldConfig fieldConfig = JsonUtils.stringToObject(
+ "{\"name\":\"event\",\"indexes\":{\"open_struct\":{}}}",
FieldConfig.class);
+ TableConfig tableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(TABLE_NAME)
+ .setFieldConfigList(List.of(fieldConfig)).build();
+ IndexLoadingConfig base = new IndexLoadingConfig(tableConfig, schema);
+ ColumnMetadata child = mock(ColumnMetadata.class);
+ when(child.getFieldSpec()).thenReturn(new DimensionFieldSpec("event$key",
DataType.INT, true));
+ SegmentMetadataImpl metadata = mock(SegmentMetadataImpl.class);
+ when(metadata.getColumnMetadataMap()).thenReturn(new
TreeMap<>(Map.of("event$key", child)));
+
+ IndexLoadingConfig derived = base.withOpenStructChildConfigs(metadata);
+ assertNotSame(derived, base);
+ assertNotNull(derived.getFieldIndexConfig("event$key"));
+ assertNull(base.getFieldIndexConfig("event$key"));
+
+ TableConfig schemaLessTable = new
TableConfigBuilder(TableType.OFFLINE).setTableName(TABLE_NAME).build();
+ IndexLoadingConfig schemaLess = new IndexLoadingConfig(schemaLessTable,
null);
+ derived = schemaLess.withKnownColumns(Set.of("segmentColumn"));
+ assertNotNull(derived.getFieldIndexConfig("segmentColumn"));
+ assertNull(schemaLess.getFieldIndexConfig("segmentColumn"));
}
@Test
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]