This is an automated email from the ASF dual-hosted git repository.
raghavyadav01 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 d5542e22727 Make an OPEN_STRUCT key behave like an ordinary column:
apply its index settings on reload, and keep its type across tiers (#19698)
d5542e22727 is described below
commit d5542e2272712f5e6a615b7cc47186235bb66900
Author: RAGHVENDRA KUMAR YADAV <[email protected]>
AuthorDate: Wed Sep 30 10:21:18 2026 -0700
Make an OPEN_STRUCT key behave like an ordinary column: apply its index
settings on reload, and keep its type across tiers (#19698)
* Apply an OPEN_STRUCT key's index settings on reload
A key's valueFieldConfigs only ever took effect while the segment was being
written: OpenStructColumnSplitter builds the per-key indexes as it writes
the
child columns, and OpenStructIndexType.createIndexHandler returned NoOp, so
nothing revisited them afterwards. Adding a range index to a key and
reloading
did nothing, while the same change on an ordinary column has always worked.
The cause is that a materialized child is a column of the segment but not
of the
table schema. FieldIndexConfigsUtil.createIndexConfigsByColName walks
schema.getColumnNames(), so a child never gets an entry, and every standard
handler derives its column set from that map -- leaving the child invisible
to
all of them.
Give OPEN_STRUCT a real handler. It derives each child's FieldIndexConfigs
from
its valueFieldConfigs entry through FieldIndexConfigsUtil.fromFieldConfig,
which
exists for exactly this case, then delegates to the ordinary per-index
handlers
so keys are indexed by the same code as every other column.
Two details the delegates need:
- The forward handler runs first, then reloadMetadata, then the rest,
mirroring
SegmentPreProcessor's own scheduling. A key written raw and later given an
inverted index needs the raw-to-dictionary conversion to land first.
- The delegates get a schema copy that includes the children.
ForwardIndexHandler
skips a column it cannot find in the schema, which is how a key silently
missed
the conversion above. The table's schema is untouched.
The dense/sparse split is not revisited; which keys are materialized is
still
decided when the segment is written.
* Keep an OPEN_STRUCT key's value type the same on both storage tiers
A materialized key carries its resolved type in its own column metadata. A
sparse key lives in a shared JSON blob that carries none, so it read back as
STRING no matter what the build inferred for it -- the same key reported
INT on
a segment that materialized it and STRING on one that did not, and a query
fanning out over both saw two types for one column.
The shape (single- vs multi-value) was already recorded in the parent
column's
metadata for exactly this reason. Record the type the same way, in a
sparseKeyTypes manifest written beside it, and resolve an undeclared sparse
key's field spec from it. A segment built before the manifest existed
carries
none and still falls back to STRING.
Also stop evaluating OpenStructTypeInference.asMultiValue twice on an
undeclared
key's first sighting: the one call both decides the shape and supplies the
elements, and it allocates.
* Apply a key's index settings through the schema, not a delegating handler
Review caught that the per-key handler resolved each child's configs from
scratch and so dropped the default-on inverted overlay that
IndexLoadingConfig.addOpenStructChildConfigs already applies. Following that
thread showed the delegation itself was wrong, not just its inputs.
A standard handler treats its config map as the whole picture for the
segment:
InvertedIndexHandler removes an inverted index from every column in the
segment
that its map does not ask for. Handing one a map that holds only OPEN_STRUCT
children therefore stripped the inverted indexes off ordinary columns. A new
test, testReloadKeepsIndexesOnOrdinaryColumns, fails that way on the
previous
commit.
None of the delegation was needed. The children are already in the config
map,
with the inverted overlay, before any handler is built. The one thing
missing
was schema visibility: ForwardIndexHandler skips a column it cannot find in
the
schema, so a key never got the raw-to-dictionary conversion an inverted
index
depends on, and every dictionary-based handler then found nothing to build
on.
So the handler goes away and SegmentPreProcessor gives the handlers a schema
copy that includes the materialized children it has configs for. The table's
own schema is untouched. needProcess() resolves the child configs the same
way
process() does -- without them a key's index settings are invisible and the
reload that would apply them is never triggered.
Verified by disabling only the schema augmentation: the regression test
fails
again with "a reload must apply the key's index settings".
* Keep a key's forward-index encoding in step with its column encoding
The end-to-end ingestion test fails on this branch: every segment carrying a
RAW OPEN_STRUCT key dies on load with
IllegalStateException: Dictionary should still exist after rebuilding
dict-encoded forward index for column: metrics$host
FieldIndexConfigsUtil.fromFieldConfig derives the dictionary entry from the
FieldConfig's encoding -- RAW disables it -- but leaves the forward entry at
ForwardIndexConfig's own JSON default, which is DICTIONARY regardless. A raw
key therefore asks for a dict-encoded forward index over a dictionary that
the
same config just disabled. ForwardIndexHandler reads that as
ENABLE_DICT_FORWARD_INDEX, rebuilds, finds no dictionary, and throws.
The inconsistency has been there since the helper was written, harmless only
because nothing acted on it: the splitter overrides both entries from its
own
dictionary decision, and before the previous commit no handler could see a
child column at all.
So derive the forward entry's encoding from the same FieldConfig the
dictionary entry comes from, preserving any explicit per-key forward
settings
through ForwardIndexConfig.Builder(other, encodingType).
testReloadOfAMixedPerKeyIndexMatrix reproduces the crash in under a second
with the integration test's own matrix -- one key inverted, one
dictionary-only, one raw -- and pins that the raw key stays raw across a
reload.
---
.../immutable/ImmutableSegmentImpl.java | 5 +-
.../impl/openstruct/OpenStructColumnSplitter.java | 45 +++-
.../segment/index/loader/SegmentPreProcessor.java | 53 +++-
.../openstruct/ImmutableOpenStructDataSource.java | 31 ++-
.../ImmutableOpenStructDataSourceTest.java | 48 ++--
.../OpenStructPerKeyIndexReloadTest.java | 281 +++++++++++++++++++++
.../OpenStructSparseDenseParityTest.java | 11 +-
.../org/apache/pinot/segment/spi/V1Constants.java | 7 +
.../segment/spi/index/FieldIndexConfigsUtil.java | 38 ++-
.../spi/index/metadata/ColumnMetadataImpl.java | 39 ++-
.../metadata/ColumnMetadataSparseKeysTest.java | 29 +++
11 files changed, 529 insertions(+), 58 deletions(-)
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentImpl.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentImpl.java
index 826f8d3663a..fbdd9a5ebc6 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentImpl.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentImpl.java
@@ -66,6 +66,7 @@ import org.apache.pinot.segment.spi.store.SegmentDirectory;
import org.apache.pinot.segment.spi.store.SegmentDirectoryPaths;
import org.apache.pinot.spi.data.ComplexFieldSpec;
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.data.readers.GenericRow;
@@ -157,14 +158,16 @@ public class ImmutableSegmentImpl implements
ImmutableSegment {
ColumnMetadata parentMetadata =
segmentMetadata.getColumnMetadataMap().get(parent);
List<String> sparseKeys = null;
Map<String, Integer> sparseMultiValueKeys = null;
+ Map<String, DataType> sparseKeyTypes = null;
if (parentMetadata instanceof ColumnMetadataImpl impl) {
sparseKeys = impl.getSparseKeys();
sparseMultiValueKeys = impl.getSparseMultiValueKeys();
+ sparseKeyTypes = impl.getSparseKeyTypes();
}
_dataSources.put(parent, new
ImmutableOpenStructDataSource((ComplexFieldSpec) fieldSpec,
openStructDenseChildren.getOrDefault(parent, Map.of()),
openStructSparseChildren.get(parent),
segmentMetadata.getTotalDocs(), sparseKeys,
- sparseMultiValueKeys));
+ sparseMultiValueKeys, sparseKeyTypes));
}
}
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/creator/impl/openstruct/OpenStructColumnSplitter.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/creator/impl/openstruct/OpenStructColumnSplitter.java
index 79664a04f22..e0efc2aa0a0 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/creator/impl/openstruct/OpenStructColumnSplitter.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/creator/impl/openstruct/OpenStructColumnSplitter.java
@@ -238,18 +238,26 @@ public class OpenStructColumnSplitter implements
ColumnarOpenStructIndexCreator
// column other documents already wrote to.
boolean multiValueKey;
boolean firstSighting = false;
+ // Held onto rather than recomputed: on an undeclared key's first sighting
the same call both decides the
+ // shape and supplies the elements, and asMultiValue allocates an array,
so evaluating it twice allocates
+ // twice for every new key.
+ Object[] inferredElements = null;
if (_presenceBitmaps.containsKey(key)) {
multiValueKey = _multiValueKeys.contains(key);
} else {
// A declaration decides the shape in both directions -- a key declared
single-value stays single-value
// even when its values are collections, because the declaration is what
the user asked for. Only an
// undeclared key takes its shape from the data.
- multiValueKey = keySpec != null
- ? !keySpec.isSingleValueField()
- : OpenStructTypeInference.asMultiValue(rawValue) != null;
+ if (keySpec != null) {
+ multiValueKey = !keySpec.isSingleValueField();
+ } else {
+ inferredElements = OpenStructTypeInference.asMultiValue(rawValue);
+ multiValueKey = inferredElements != null;
+ }
firstSighting = true;
}
- Object[] elements = multiValueKey ?
OpenStructTypeInference.asMultiValue(rawValue) : null;
+ Object[] elements = !multiValueKey ? null
+ : inferredElements != null ? inferredElements :
OpenStructTypeInference.asMultiValue(rawValue);
if (elements != null && elements.length == 0) {
// No elements, so no value and nothing to infer a type from. A
materialized multi-value column has no
// empty state, so treating this as present would mean inventing one;
the key is simply not in this
@@ -493,13 +501,19 @@ public class OpenStructColumnSplitter implements
ColumnarOpenStructIndexCreator
keySpec.getDefaultNullValue());
}
+ /// The type a key resolves to, whichever tier it ends up on: a declaration
when there is one, else what the
+ /// values inferred. This is the single source both [#writeDenseKeyColumn]
and the sparse type manifest read, so
+ /// the two tiers cannot report different types for the same key.
+ private DataType resolvedKeyType(String key) {
+ FieldSpec keySpec = _childFieldSpecs.get(key);
+ return keySpec != null ? keySpec.getDataType() :
_inferredTypes.getOrDefault(key, DataType.STRING);
+ }
+
private void writeDenseKeyColumn(String key)
throws IOException {
String materializedCol =
OpenStructNaming.materializedColumnName(_columnName, key);
FieldSpec keySpec = _childFieldSpecs.get(key);
- DataType valueType = keySpec != null
- ? keySpec.getDataType()
- : _inferredTypes.getOrDefault(key, DataType.STRING);
+ DataType valueType = resolvedKeyType(key);
DataType storedType = valueType.getStoredType();
RoaringBitmap presence = _presenceBitmaps.get(key);
List<Object> values = _values.get(key);
@@ -849,6 +863,23 @@ public class OpenStructColumnSplitter implements
ColumnarOpenStructIndexCreator
throw new RuntimeException("Failed to serialize sparse multi-value
key manifest", e);
}
}
+ // Type, for the same reason as shape: the dense/sparse split is a
tuning decision and must not change how a
+ // key reads. A dense key carries its type in its own column metadata;
the blob a sparse key lives in is JSON
+ // and carries none, so without this the same key read as INT on a
segment that materialized it and STRING on
+ // one that did not, and a query fanning out over both saw two types for
one column.
+ Map<String, DataType> sparseKeyTypes = new LinkedHashMap<>();
+ for (String key : sparseKeys) {
+ sparseKeyTypes.put(key, resolvedKeyType(key).getStoredType());
+ }
+ if (!sparseKeyTypes.isEmpty()) {
+ try {
+
props.setProperty(V1Constants.MetadataKeys.Column.getKeyFor(_columnName,
+ V1Constants.MetadataKeys.Column.SPARSE_KEY_TYPES),
+ JsonUtils.objectToString(sparseKeyTypes));
+ } catch (IOException e) {
+ throw new RuntimeException("Failed to serialize sparse key type
manifest", e);
+ }
+ }
}
_materializedColumnMetadata.put(_columnName, props);
}
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/loader/SegmentPreProcessor.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/loader/SegmentPreProcessor.java
index f8dbc0dddce..1bac733774a 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/loader/SegmentPreProcessor.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/loader/SegmentPreProcessor.java
@@ -24,6 +24,7 @@ import java.io.IOException;
import java.util.ArrayList;
import java.util.LinkedHashSet;
import java.util.List;
+import java.util.Map;
import java.util.ServiceLoader;
import java.util.Set;
import javax.annotation.Nullable;
@@ -42,11 +43,14 @@ import
org.apache.pinot.segment.local.startree.StarTreeBuilderUtils;
import org.apache.pinot.segment.local.startree.v2.builder.MultipleTreesBuilder;
import
org.apache.pinot.segment.local.startree.v2.builder.StarTreeV2BuilderConfig;
import org.apache.pinot.segment.local.utils.SegmentOperationsThrottlerSet;
+import org.apache.pinot.segment.spi.ColumnMetadata;
import org.apache.pinot.segment.spi.V1Constants;
+import org.apache.pinot.segment.spi.index.FieldIndexConfigs;
import org.apache.pinot.segment.spi.index.IndexHandler;
import org.apache.pinot.segment.spi.index.IndexService;
import org.apache.pinot.segment.spi.index.IndexType;
import org.apache.pinot.segment.spi.index.StandardIndexes;
+import org.apache.pinot.segment.spi.index.metadata.ColumnMetadataImpl;
import org.apache.pinot.segment.spi.index.metadata.SegmentMetadataImpl;
import
org.apache.pinot.segment.spi.index.multicolumntext.MultiColumnTextIndexConstants;
import
org.apache.pinot.segment.spi.index.multicolumntext.MultiColumnTextMetadata;
@@ -56,6 +60,7 @@ import
org.apache.pinot.segment.spi.store.SegmentDirectoryPaths;
import org.apache.pinot.segment.spi.utils.SegmentMetadataUtils;
import org.apache.pinot.spi.config.table.MultiColumnTextIndexConfig;
import org.apache.pinot.spi.config.table.TableConfig;
+import org.apache.pinot.spi.data.FieldSpec;
import org.apache.pinot.spi.data.Schema;
import org.apache.pinot.spi.plugin.PluginManager;
import org.slf4j.Logger;
@@ -179,7 +184,8 @@ public class SegmentPreProcessor implements AutoCloseable {
// build dict-id-based indexes on top of the new shared dictionary. If
this order is violated, downstream
// handlers fail with an IllegalStateException because the dictionary
they require does not yet exist.
// Any future change to handler scheduling MUST preserve:
ForwardIndexHandler → reloadMetadata → other handlers.
- IndexHandler forwardHandler = createHandler(StandardIndexes.forward());
+ Map<String, FieldIndexConfigs> configsByCol =
_indexLoadingConfig.getFieldIndexConfigByColName();
+ IndexHandler forwardHandler = createHandler(StandardIndexes.forward(),
configsByCol);
indexHandlers.add(forwardHandler);
forwardHandler.updateIndices(segmentWriter);
_segmentDirectory.reloadMetadata();
@@ -187,7 +193,7 @@ public class SegmentPreProcessor implements AutoCloseable {
// Now that ForwardIndexHandler.updateIndices has been updated, we can
run all other indexes in any order
for (IndexType<?, ?, ?> type :
IndexService.getInstance().getAllIndexes()) {
if (type != StandardIndexes.forward()) {
- IndexHandler handler = createHandler(type);
+ IndexHandler handler = createHandler(type, configsByCol);
indexHandlers.add(handler);
handler.updateIndices(segmentWriter);
}
@@ -230,11 +236,41 @@ public class SegmentPreProcessor implements AutoCloseable
{
}
}
- private IndexHandler createHandler(IndexType<?, ?, ?> type) {
- return type.createIndexHandler(_segmentDirectory,
_indexLoadingConfig.getFieldIndexConfigByColName(), _schema,
+ private IndexHandler createHandler(IndexType<?, ?, ?> type, Map<String,
FieldIndexConfigs> configsByCol) {
+ return type.createIndexHandler(_segmentDirectory, configsByCol,
schemaWithMaterializedChildColumns(configsByCol),
_tableConfig);
}
+ /// The table schema plus one entry per OPEN_STRUCT materialized child that
has index configs.
+ ///
+ /// A materialized child (`col$key`) is a real column of the segment but not
of the table schema, and
+ /// [ForwardIndexHandler] skips a column it cannot find in the schema. That
is why a key never got the
+ /// raw-to-dictionary conversion an inverted or range index depends on: its
configs were in the config map
+ /// (see [IndexLoadingConfig#withOpenStructChildConfigs]) but the column was
invisible to the one handler that
+ /// had to act first. The child's spec comes from segment metadata; the
table's own schema is left untouched.
+ private Schema schemaWithMaterializedChildColumns(Map<String,
FieldIndexConfigs> configsByCol) {
+ Map<String, ColumnMetadata> columnMetadataMap =
_segmentDirectory.getSegmentMetadata().getColumnMetadataMap();
+ Schema augmented = null;
+ for (Map.Entry<String, ColumnMetadata> entry :
columnMetadataMap.entrySet()) {
+ String column = entry.getKey();
+ if (_schema.hasColumn(column) || !configsByCol.containsKey(column)) {
+ continue;
+ }
+ if (!(entry.getValue() instanceof ColumnMetadataImpl columnMetadata) ||
!columnMetadata.isMaterializedChild()) {
+ continue;
+ }
+ if (augmented == null) {
+ augmented = new Schema();
+ augmented.setSchemaName(_schema.getSchemaName());
+ for (FieldSpec fieldSpec : _schema.getAllFieldSpecs()) {
+ augmented.addField(fieldSpec);
+ }
+ }
+ augmented.addField(columnMetadata.getFieldSpec());
+ }
+ return augmented == null ? _schema : augmented;
+ }
+
/// This method checks if there is any discrepancy between the segment and
current table config and schema.
/// If so, it returns true indicating the segment needs to be reprocessed.
Right now, the default columns,
/// all types of indices and column min/max values are checked against
what's set in table config and schema.
@@ -253,9 +289,14 @@ public class SegmentPreProcessor implements AutoCloseable {
LOGGER.info("Found default columns need updates in segment: {}",
segmentName);
return true;
}
- // Check if there is need to update single-column indices, like inverted
index, json index etc.
+ // Check if there is need to update single-column indices, like inverted
index, json index etc. The
+ // OPEN_STRUCT child configs are resolved here for the same reason
process() resolves them: without them a
+ // key's index settings are invisible, and a reload that would apply
them is never triggered in the first
+ // place. A derived copy is used so a config shared across segments does
not pick up this segment's children.
+ Map<String, FieldIndexConfigs> configsByCol =
+
_indexLoadingConfig.withOpenStructChildConfigs(segmentMetadata).getFieldIndexConfigByColName();
for (IndexType<?, ?, ?> type :
IndexService.getInstance().getAllIndexes()) {
- if (createHandler(type).needUpdateIndices(segmentReader)) {
+ if (createHandler(type,
configsByCol).needUpdateIndices(segmentReader)) {
LOGGER.info("Found index type: {} needs updates in segment: {}",
type, segmentName);
return true;
}
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/openstruct/ImmutableOpenStructDataSource.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/openstruct/ImmutableOpenStructDataSource.java
index 6480d55d691..3176bbf9c88 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/openstruct/ImmutableOpenStructDataSource.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/segment/index/openstruct/ImmutableOpenStructDataSource.java
@@ -58,6 +58,10 @@ public class ImmutableOpenStructDataSource extends
BaseDataSource implements Ope
/// column of their own.
@Nullable
private final Map<String, Integer> _sparseMultiValueKeys;
+ /// Sparse keys mapped to the type the segment build resolved for them, read
from the parent column's metadata
+ /// for the same reason as the shape above: the blob is JSON and carries no
type of its own.
+ @Nullable
+ private final Map<String, FieldSpec.DataType> _sparseKeyTypes;
@Nullable
private final OpenStructSparseBlobReader _sparseBlobReader;
private final ConcurrentHashMap<String, DataSource>
_sparseKeyDataSourceCache;
@@ -65,13 +69,15 @@ public class ImmutableOpenStructDataSource extends
BaseDataSource implements Ope
public ImmutableOpenStructDataSource(ComplexFieldSpec fieldSpec, Map<String,
DataSource> perKeyDataSources,
@Nullable DataSource sparseDataSource, DataSourceMetadata
dataSourceMetadata,
ColumnIndexContainer indexContainer, @Nullable List<String> sparseKeys,
- @Nullable Map<String, Integer> sparseMultiValueKeys) {
+ @Nullable Map<String, Integer> sparseMultiValueKeys,
+ @Nullable Map<String, FieldSpec.DataType> sparseKeyTypes) {
super(dataSourceMetadata, indexContainer);
_fieldSpec = fieldSpec;
_perKeyDataSources = perKeyDataSources;
_sparseDataSource = sparseDataSource;
_sparseKeys = sparseKeys != null ? Set.copyOf(sparseKeys) : null;
_sparseMultiValueKeys = sparseMultiValueKeys != null ?
Map.copyOf(sparseMultiValueKeys) : null;
+ _sparseKeyTypes = sparseKeyTypes != null ? Map.copyOf(sparseKeyTypes) :
null;
if (sparseDataSource != null) {
ForwardIndexReader<?> blobFwd = sparseDataSource.getForwardIndex();
_sparseBlobReader = blobFwd != null
@@ -93,10 +99,11 @@ public class ImmutableOpenStructDataSource extends
BaseDataSource implements Ope
/// (`SELECT open_struct_col`) is handled by the query layer, not the
storage layer.
public ImmutableOpenStructDataSource(ComplexFieldSpec fieldSpec, Map<String,
DataSource> perKeyDataSources,
@Nullable DataSource sparseDataSource, int numDocs, @Nullable
List<String> sparseKeys,
- @Nullable Map<String, Integer> sparseMultiValueKeys) {
+ @Nullable Map<String, Integer> sparseMultiValueKeys,
+ @Nullable Map<String, FieldSpec.DataType> sparseKeyTypes) {
this(fieldSpec, perKeyDataSources, sparseDataSource,
new ImmutableOpenStructDataSourceMetadata(fieldSpec, numDocs),
- new ColumnIndexContainer.FromMap.Builder().build(), sparseKeys,
sparseMultiValueKeys);
+ new ColumnIndexContainer.FromMap.Builder().build(), sparseKeys,
sparseMultiValueKeys, sparseKeyTypes);
}
@Override
@@ -118,11 +125,14 @@ public class ImmutableOpenStructDataSource extends
BaseDataSource implements Ope
k -> new SparseKeyDataSource(getValueFieldSpec(k), _sparseBlobReader,
maxNumValues(k)));
}
- /// Field spec for a key's values, with an undeclared sparse key's shape
taken from the segment's sparse
- /// multi-value manifest. Which tier a key lands on is a tuning decision, so
it must not decide the key's
- /// shape: without this, the same rows would report `STRING[]` on a segment
that materialized the key and a
- /// scalar `STRING` holding `["a","b"]` on one that put it in the blob, and
a query fanning out over both
- /// would see two shapes for one column.
+ /// Field spec for a key's values, with an undeclared sparse key's shape and
type taken from the segment's
+ /// sparse manifests. Which tier a key lands on is a tuning decision, so it
must not decide either: without
+ /// this, the same rows would report `INT` on a segment that materialized
the key and `STRING` on one that put
+ /// it in the blob -- and `STRING[]` versus a scalar `STRING` holding
`["a","b"]` for shape -- so a query
+ /// fanning out over both would see two types for one column.
+ ///
+ /// A segment built before the manifests existed lists nothing, and an
unlisted key falls back to the
+ /// single-value STRING every sparse key used to read as.
@Override
public FieldSpec getValueFieldSpec(String key) {
FieldSpec childFieldSpec = _fieldSpec.getChildFieldSpec(key);
@@ -130,7 +140,10 @@ public class ImmutableOpenStructDataSource extends
BaseDataSource implements Ope
return childFieldSpec;
}
boolean singleValue = _sparseMultiValueKeys == null ||
!_sparseMultiValueKeys.containsKey(key);
- return new DimensionFieldSpec(key, FieldSpec.DataType.STRING, singleValue);
+ FieldSpec.DataType dataType = _sparseKeyTypes != null
+ ? _sparseKeyTypes.getOrDefault(key, FieldSpec.DataType.STRING)
+ : FieldSpec.DataType.STRING;
+ return new DimensionFieldSpec(key, dataType, singleValue);
}
/// Longest value a multi-value sparse key holds, which is what the readers
over it size their buffers from.
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/ImmutableOpenStructDataSourceTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/ImmutableOpenStructDataSourceTest.java
index 383b36b32d2..8af04da884c 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/ImmutableOpenStructDataSourceTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/ImmutableOpenStructDataSourceTest.java
@@ -94,7 +94,7 @@ public class ImmutableOpenStructDataSourceTest {
sparseDs,
meta,
container,
- null, null);
+ null, null, null);
assertSame(ds.getDataSource("clicks"), clicksDs);
// Absent key with mock sparse (no real forward index) resolves to an
all-null STRING source
@@ -109,7 +109,7 @@ public class ImmutableOpenStructDataSourceTest {
public void testAbsentDeclaredKeyReadsDeclaredDefault() {
ComplexFieldSpec spec = new ComplexFieldSpec("event",
DataType.OPEN_STRUCT, true,
Map.of("score", new DimensionFieldSpec("score", DataType.STRING, true,
"N/A")));
- ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(spec,
Map.of(), null, 7, null, null);
+ ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(spec,
Map.of(), null, 7, null, null, null);
DataSource scoreDs = ds.getDataSource("score");
assertTrue(scoreDs instanceof NullDataSource);
@@ -131,7 +131,7 @@ public class ImmutableOpenStructDataSourceTest {
null,
meta,
container,
- null, null);
+ null, null, null);
assertTrue(ds.isMaterialized("clicks"));
assertFalse(ds.isMaterialized("absent"));
@@ -148,7 +148,7 @@ public class ImmutableOpenStructDataSourceTest {
null,
meta,
container,
- null, null);
+ null, null, null);
assertTrue(ds.isFullyMaterialized());
}
@@ -165,7 +165,7 @@ public class ImmutableOpenStructDataSourceTest {
sparseDs,
meta,
container,
- null, null);
+ null, null, null);
assertFalse(ds.isFullyMaterialized());
}
@@ -181,7 +181,7 @@ public class ImmutableOpenStructDataSourceTest {
null,
meta,
container,
- null, null);
+ null, null, null);
ComplexFieldSpec fieldSpec = ds.getFieldSpec();
assertNotNull(fieldSpec);
@@ -203,7 +203,7 @@ public class ImmutableOpenStructDataSourceTest {
null,
topMeta,
container,
- null, null);
+ null, null, null);
assertSame(ds.getDataSourceMetadata("clicks"), clicksMeta);
assertEquals(ds.getDataSourceMetadata("absent").getDataType(),
DataType.STRING);
@@ -222,7 +222,7 @@ public class ImmutableOpenStructDataSourceTest {
null,
meta,
container,
- null, null);
+ null, null, null);
assertEquals(ds.getDataSources(), perKeyMap);
}
@@ -238,7 +238,7 @@ public class ImmutableOpenStructDataSourceTest {
null,
meta,
container,
- null, null);
+ null, null, null);
assertSame(ds.getDataSourceMetadata(), meta);
assertSame(ds.getIndexContainer(), container);
@@ -252,7 +252,7 @@ public class ImmutableOpenStructDataSourceTest {
Map.of("clicks", clicksDs),
null,
42,
- null, null);
+ null, null, null);
DataSourceMetadata meta = ds.getDataSourceMetadata();
assertNotNull(meta);
@@ -332,7 +332,7 @@ public class ImmutableOpenStructDataSourceTest {
Map.of(),
sparseDs,
blobs.length,
- List.of("region"), null);
+ List.of("region"), null, null);
DataSource regionDs = ds.getDataSource("region");
assertNotNull(regionDs);
@@ -354,7 +354,7 @@ public class ImmutableOpenStructDataSourceTest {
Map.of(),
sparseDs,
blobs.length,
- null, null);
+ null, null, null);
DataSource anyDs = ds.getDataSource("anything");
assertNotNull(anyDs);
@@ -374,7 +374,7 @@ public class ImmutableOpenStructDataSourceTest {
Map.of(),
sparseDs,
blobs.length,
- List.of("latencyMs"), null);
+ List.of("latencyMs"), null, null);
DataSource latDs = ds.getDataSource("latencyMs");
assertNotNull(latDs);
@@ -390,14 +390,14 @@ public class ImmutableOpenStructDataSourceTest {
doReturn(mockJsonIdx).when(sparseDs).getJsonIndex();
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), Map.of(), sparseDs, blobs.length, null, null);
+ openStructSpec("event"), Map.of(), sparseDs, blobs.length, null, null,
null);
assertSame(ds.getSparseJsonIndex(), mockJsonIdx);
}
@Test
public void testGetSparseJsonIndexNullWhenNoSparseColumn() {
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), Map.of(), null, 5, null, null);
+ openStructSpec("event"), Map.of(), null, 5, null, null, null);
assertNull(ds.getSparseJsonIndex());
}
@@ -411,7 +411,7 @@ public class ImmutableOpenStructDataSourceTest {
perKey.put("name", nameDs);
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), perKey, null, 2, null, null);
+ openStructSpec("event"), perKey, null, 2, null, null, null);
Map<String, Object> doc0 = ds.getMapValue(0);
assertNotNull(doc0);
@@ -441,7 +441,7 @@ public class ImmutableOpenStructDataSourceTest {
DataSource rawDs = mockRawDenseDataSource(storedType, expected);
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), Map.of("k", rawDs), null, 2, null, null);
+ openStructSpec("event"), Map.of("k", rawDs), null, 2, null, null,
null);
Map<String, Object> doc0 = ds.getMapValue(0);
assertNotNull(doc0);
@@ -505,7 +505,7 @@ public class ImmutableOpenStructDataSourceTest {
DataSource clicksDs = mockDenseDataSource(DataType.INT, 10, true);
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), Map.of("clicks", clicksDs), null, 2, null,
null);
+ openStructSpec("event"), Map.of("clicks", clicksDs), null, 2, null,
null, null);
// doc 1 has null for clicks → no keys → null map
Map<String, Object> doc1 = ds.getMapValue(1);
@@ -518,7 +518,7 @@ public class ImmutableOpenStructDataSourceTest {
DataSource sparseDs = mockSparseDataSource("{\"rare_key\":\"val\"}", true);
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), Map.of("clicks", clicksDs), sparseDs, 2,
null, null);
+ openStructSpec("event"), Map.of("clicks", clicksDs), sparseDs, 2,
null, null, null);
Map<String, Object> doc0 = ds.getMapValue(0);
assertNotNull(doc0);
@@ -538,7 +538,7 @@ public class ImmutableOpenStructDataSourceTest {
DataSource sparseDs = mockSparseDataSource("{\"region\":\"us\"}", false);
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), Map.of("event", eventKeyDs), sparseDs, 2,
null, null);
+ openStructSpec("event"), Map.of("event", eventKeyDs), sparseDs, 2,
null, null, null);
try (MapValueReader reader = ds.openMapValueReader()) {
Map<String, Object> doc0 = reader.getMapValue(0);
@@ -553,7 +553,7 @@ public class ImmutableOpenStructDataSourceTest {
DataSource sparseDs = mockSparseDataSource("not-json", false);
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), Map.of(), sparseDs, 2, null, null);
+ openStructSpec("event"), Map.of(), sparseDs, 2, null, null, null);
RuntimeException e = expectThrows(RuntimeException.class, () ->
ds.getMapValue(0));
assertTrue(e.getMessage().contains("docId 0"));
@@ -574,7 +574,7 @@ public class ImmutableOpenStructDataSourceTest {
perKey.put("b", ds2);
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), perKey, null, 1, null, null);
+ openStructSpec("event"), perKey, null, 1, null, null, null);
MapValueReader reader = ds.openMapValueReader();
reader.getMapValue(0);
@@ -609,7 +609,7 @@ public class ImmutableOpenStructDataSourceTest {
DataSource sparseDs = mockSparseDataSource("{\"rare_key\":\"val\"}", true);
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), Map.of(), sparseDs, 2, null, null);
+ openStructSpec("event"), Map.of(), sparseDs, 2, null, null, null);
// doc 1: sparse is null
Map<String, Object> doc1 = ds.getMapValue(1);
@@ -619,7 +619,7 @@ public class ImmutableOpenStructDataSourceTest {
@Test
public void testGetMapValueEmptySegment() {
ImmutableOpenStructDataSource ds = new ImmutableOpenStructDataSource(
- openStructSpec("event"), Map.of(), null, 0, null, null);
+ openStructSpec("event"), Map.of(), null, 0, null, null, null);
assertNull(ds.getMapValue(0));
}
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/OpenStructPerKeyIndexReloadTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/OpenStructPerKeyIndexReloadTest.java
new file mode 100644
index 00000000000..54d710a237c
--- /dev/null
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/OpenStructPerKeyIndexReloadTest.java
@@ -0,0 +1,281 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.segment.local.segment.index.openstruct;
+
+import com.fasterxml.jackson.databind.node.ObjectNode;
+import java.io.File;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import org.apache.commons.io.FileUtils;
+import
org.apache.pinot.segment.local.segment.creator.impl.SegmentIndexCreationDriverImpl;
+import org.apache.pinot.segment.local.segment.index.loader.IndexLoadingConfig;
+import org.apache.pinot.segment.local.segment.index.loader.SegmentPreProcessor;
+import org.apache.pinot.segment.local.segment.readers.GenericRowRecordReader;
+import org.apache.pinot.segment.local.segment.store.SegmentLocalFSDirectory;
+import org.apache.pinot.segment.spi.ColumnMetadata;
+import org.apache.pinot.segment.spi.V1Constants;
+import org.apache.pinot.segment.spi.creator.SegmentGeneratorConfig;
+import org.apache.pinot.segment.spi.index.IndexType;
+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.spi.config.table.FieldConfig;
+import org.apache.pinot.spi.config.table.OpenStructIndexConfig;
+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.OpenStructNaming;
+import org.apache.pinot.spi.data.Schema;
+import org.apache.pinot.spi.data.readers.GenericRow;
+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.AfterClass;
+import org.testng.annotations.BeforeClass;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertTrue;
+
+
+/// A key's index settings must be applied by a reload, not only when the
segment was written.
+///
+/// Per-key indexes used to be built solely by `OpenStructColumnSplitter`
while it wrote the child columns, so they
+/// were frozen at creation: adding a range index to a key and reloading did
nothing, while the same change on an
+/// ordinary column has always worked.
+public class OpenStructPerKeyIndexReloadTest {
+ private static final File TMP_DIR =
+ new File(FileUtils.getTempDirectory(),
OpenStructPerKeyIndexReloadTest.class.getSimpleName());
+ private static final String TABLE = "openStructPerKeyIndex";
+ private static final String COLUMN = "props";
+ private static final String KEY = "region";
+ private static final String CHILD =
OpenStructNaming.materializedColumnName(COLUMN, KEY);
+ private static final String ORDINARY_COLUMN = "id";
+ private static final int NUM_DOCS = 64;
+
+ @BeforeClass
+ public void setUp()
+ throws Exception {
+ FileUtils.deleteQuietly(TMP_DIR);
+ FileUtils.forceMkdir(TMP_DIR);
+ }
+
+ @AfterClass
+ public void tearDown() {
+ FileUtils.deleteQuietly(TMP_DIR);
+ }
+
+ /// The regression: build with the key un-indexed, then reload with an
inverted index configured for it.
+ @Test
+ public void testReloadBuildsAPerKeyIndexThatWasNotThereAtCreation()
+ throws Exception {
+ File segmentDir = buildSegment("reloadAddsIndex", withoutPerKeyIndex());
+
+ assertFalse(hasIndex(segmentDir, CHILD, StandardIndexes.inverted()),
+ "precondition: the key must start without an inverted index");
+
+ reload(segmentDir, withPerKeyInvertedIndex());
+
+ assertTrue(hasIndex(segmentDir, CHILD, StandardIndexes.inverted()),
+ "a reload must apply the key's index settings, the way it does for any
other column");
+ }
+
+ /// A segment whose keys already carry what the config asks for must need no
reprocessing, and a reload of one
+ /// must leave the indexes alone -- otherwise every reload rewrites indexes
that are already correct.
+ @Test
+ public void testReloadIsANoOpWhenTheKeyIndexIsAlreadyPresent()
+ throws Exception {
+ File segmentDir = buildSegment("reloadNoOp", withPerKeyInvertedIndex());
+ assertTrue(hasIndex(segmentDir, CHILD, StandardIndexes.inverted()),
+ "precondition: creation already built the index");
+
+ try (SegmentDirectory directory = openDirectory(segmentDir);
+ SegmentPreProcessor preProcessor =
+ new SegmentPreProcessor(directory,
indexLoadingConfig(withPerKeyInvertedIndex()))) {
+ assertFalse(preProcessor.needProcess(), "a key whose index already
matches needs no reprocessing");
+ }
+
+ reload(segmentDir, withPerKeyInvertedIndex());
+ assertTrue(hasIndex(segmentDir, CHILD, StandardIndexes.inverted()), "and
the index must survive a reload");
+ }
+
+ /// The sparse blob holds every unmaterialized key in one column, so a
per-key setting cannot mean anything for
+ /// it. It must be left alone rather than indexed as if it were a key.
+ @Test
+ public void testSparseBlobColumnIsNotTreatedAsAKey()
+ throws Exception {
+ File segmentDir = buildSegment("sparseUntouched",
withPerKeyInvertedIndex());
+ String sparse = OpenStructNaming.sparseColumnName(COLUMN);
+ reload(segmentDir, withPerKeyInvertedIndex());
+
+ SegmentMetadataImpl metadata = new SegmentMetadataImpl(segmentDir);
+ if (metadata.getColumnMetadataMap().containsKey(sparse)) {
+ assertFalse(hasIndex(segmentDir, sparse, StandardIndexes.inverted()),
+ "the shared blob column must not get a per-key inverted index");
+ }
+ }
+
+ /// An ordinary column's indexes must survive a reload of a table that also
has OPEN_STRUCT keys.
+ @Test
+ public void testReloadKeepsIndexesOnOrdinaryColumns()
+ throws Exception {
+ File segmentDir = buildSegment("ordinaryUntouched",
withPerKeyInvertedIndex());
+ assertTrue(hasIndex(segmentDir, ORDINARY_COLUMN,
StandardIndexes.inverted()),
+ "precondition: creation built the ordinary column's inverted index");
+
+ reload(segmentDir, withPerKeyInvertedIndex());
+
+ assertTrue(hasIndex(segmentDir, ORDINARY_COLUMN,
StandardIndexes.inverted()),
+ "an ordinary column's inverted index must survive the reload");
+ }
+
+ /// The per-key index matrix the end-to-end ingestion test uses: one key
inverted, one dictionary-only, one raw.
+ /// A reload must leave all three as configured -- a raw key in particular
must not be dragged into a dictionary.
+ @Test
+ public void testReloadOfAMixedPerKeyIndexMatrix()
+ throws Exception {
+ File segmentDir = buildSegment("mixedMatrix", withMixedMatrix());
+ reload(segmentDir, withMixedMatrix());
+
+ String views = OpenStructNaming.materializedColumnName(COLUMN, "views");
+ String cpu = OpenStructNaming.materializedColumnName(COLUMN, "cpu");
+ String host = OpenStructNaming.materializedColumnName(COLUMN, "host");
+ assertTrue(hasIndex(segmentDir, views, StandardIndexes.inverted()), "the
inverted key keeps its index");
+ assertTrue(hasIndex(segmentDir, cpu, StandardIndexes.dictionary()), "the
dictionary key keeps its dictionary");
+ assertFalse(hasIndex(segmentDir, host, StandardIndexes.dictionary()), "the
raw key must stay raw");
+ }
+
+ private static OpenStructIndexConfig withMixedMatrix() {
+ FieldConfig views = new FieldConfig.Builder("views")
+ .withIndexes(JsonUtils.objectToJsonNode(Map.of("inverted", Map.of())))
+ .build();
+ FieldConfig cpu = new FieldConfig.Builder("cpu")
+ .withEncodingType(FieldConfig.EncodingType.DICTIONARY)
+ .build();
+ FieldConfig host = new FieldConfig.Builder("host")
+ .withEncodingType(FieldConfig.EncodingType.RAW)
+ .build();
+ return new OpenStructIndexConfig(false, null, 3, Set.of("views", "cpu",
"host"), 0.5,
+ List.of(views, cpu, host));
+ }
+
+ // ---------------------------------------------------------------- fixtures
+
+ private static OpenStructIndexConfig withoutPerKeyIndex() {
+ // Every key raw and un-indexed, so the child starts with nothing for the
reload to find.
+ FieldConfig raw = new
FieldConfig.Builder("default").withEncodingType(FieldConfig.EncodingType.RAW).build();
+ // First arg is `disabled`, not `enabled`.
+ return new OpenStructIndexConfig(false, raw, -1, null, 0.0, null);
+ }
+
+ private static OpenStructIndexConfig withPerKeyInvertedIndex() {
+ FieldConfig raw = new
FieldConfig.Builder("default").withEncodingType(FieldConfig.EncodingType.RAW).build();
+ ObjectNode inverted = JsonUtils.newObjectNode();
+ inverted.set("inverted", JsonUtils.newObjectNode());
+ FieldConfig keyConfig = new FieldConfig.Builder(KEY)
+ .withEncodingType(FieldConfig.EncodingType.DICTIONARY)
+ .withIndexes(inverted)
+ .build();
+ return new OpenStructIndexConfig(false, raw, -1, Set.of(KEY), 0.0,
List.of(keyConfig));
+ }
+
+ private static Schema schema() {
+ return new Schema.SchemaBuilder().setSchemaName(TABLE)
+ .addField(new ComplexFieldSpec(COLUMN, FieldSpec.DataType.OPEN_STRUCT,
true, Map.of(
+ "views", new DimensionFieldSpec("views", FieldSpec.DataType.LONG,
true),
+ "cpu", new DimensionFieldSpec("cpu", FieldSpec.DataType.DOUBLE,
true),
+ "host", new DimensionFieldSpec("host", FieldSpec.DataType.STRING,
true))))
+ .addSingleValueDimension(ORDINARY_COLUMN, FieldSpec.DataType.STRING)
+ .build();
+ }
+
+ private static TableConfig tableConfig(OpenStructIndexConfig osConfig) {
+ ObjectNode indexes = JsonUtils.newObjectNode();
+ indexes.set("open_struct", JsonUtils.objectToJsonNode(osConfig));
+ return new TableConfigBuilder(TableType.OFFLINE).setTableName(TABLE)
+ .setFieldConfigList(List.of(new
FieldConfig.Builder(COLUMN).withIndexes(indexes).build()))
+ .setInvertedIndexColumns(List.of(ORDINARY_COLUMN))
+ .setNullHandlingEnabled(true)
+ .build();
+ }
+
+ private static File buildSegment(String segmentName, OpenStructIndexConfig
osConfig)
+ throws Exception {
+ List<GenericRow> rows = new ArrayList<>(NUM_DOCS);
+ for (int docId = 0; docId < NUM_DOCS; docId++) {
+ Map<String, Object> props = new HashMap<>();
+ props.put(KEY, docId % 2 == 0 ? "us" : "eu");
+ props.put("views", (long) docId);
+ props.put("cpu", docId * 0.5);
+ props.put("host", "host-" + (docId % 5));
+ GenericRow row = new GenericRow();
+ row.putValue(COLUMN, props);
+ row.putValue(ORDINARY_COLUMN, "id-" + (docId % 8));
+ rows.add(row);
+ }
+ SegmentGeneratorConfig config = new
SegmentGeneratorConfig(tableConfig(osConfig), schema());
+ config.setOutDir(new File(TMP_DIR, segmentName).getAbsolutePath());
+ config.setSegmentName(segmentName);
+ SegmentIndexCreationDriverImpl driver = new
SegmentIndexCreationDriverImpl();
+ driver.init(config, new GenericRowRecordReader(rows));
+ driver.build();
+ return driver.getOutputDirectory();
+ }
+
+ private static IndexLoadingConfig indexLoadingConfig(OpenStructIndexConfig
osConfig) {
+ return new IndexLoadingConfig(tableConfig(osConfig), schema());
+ }
+
+ private static SegmentDirectory openDirectory(File segmentDir)
+ throws Exception {
+ return new SegmentLocalFSDirectory(segmentDir, new
SegmentMetadataImpl(segmentDir), ReadMode.mmap);
+ }
+
+ private static void reload(File segmentDir, OpenStructIndexConfig osConfig)
+ throws Exception {
+ try (SegmentDirectory directory = openDirectory(segmentDir);
+ SegmentPreProcessor preProcessor = new SegmentPreProcessor(directory,
indexLoadingConfig(osConfig))) {
+ preProcessor.process(null);
+ }
+ }
+
+ private static boolean hasIndex(File segmentDir, String column, IndexType<?,
?, ?> indexType)
+ throws Exception {
+ SegmentMetadataImpl metadata = new SegmentMetadataImpl(segmentDir);
+ ColumnMetadata columnMetadata = metadata.getColumnMetadataFor(column);
+ if (columnMetadata == null) {
+ return false;
+ }
+ try (SegmentDirectory directory =
+ new SegmentLocalFSDirectory(segmentDir, metadata, ReadMode.mmap);
+ SegmentDirectory.Reader reader = directory.createReader()) {
+ return reader.hasIndexFor(column, indexType);
+ }
+ }
+
+ static {
+ // Keep the V1 constants class loaded so segment metadata reads resolve
the same way the loader does.
+ assert V1Constants.MetadataKeys.Column.PARENT_COLUMN != null;
+ }
+}
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/OpenStructSparseDenseParityTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/OpenStructSparseDenseParityTest.java
index 4f9a7f41ac2..3e5ca1f7995 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/OpenStructSparseDenseParityTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/segment/index/openstruct/OpenStructSparseDenseParityTest.java
@@ -165,13 +165,18 @@ public class OpenStructSparseDenseParityTest {
assertForwardValuesParity(denseDs, sparseDs, "latencyMs",
FieldSpec.DataType.LONG);
assertNullBitmapParity(denseDs, sparseDs, "latencyMs");
- // Undeclared key: dense infers INT, sparse defaults to STRING.
+ // An undeclared key takes its type from the data, and the tier it lands
on must not change that: the
+ // segment records what the build inferred, so the key reads as INT on
both sides. Without the recorded
+ // type the sparse side fell back to STRING and the same key reported
two types across segments.
+ assertForwardValuesParity(denseDs, sparseDs, "freeform",
FieldSpec.DataType.INT);
+ assertNullBitmapParity(denseDs, sparseDs, "freeform");
+
DataSource denseFreeform = denseDs.getDataSource("freeform");
DataSource sparseFreeform = sparseDs.getDataSource("freeform");
assertNotNull(denseFreeform);
assertNotNull(sparseFreeform);
assertEquals(denseFreeform.getDataSourceMetadata().getDataType().getStoredType(),
FieldSpec.DataType.INT);
-
assertEquals(sparseFreeform.getDataSourceMetadata().getDataType().getStoredType(),
FieldSpec.DataType.STRING);
+
assertEquals(sparseFreeform.getDataSourceMetadata().getDataType().getStoredType(),
FieldSpec.DataType.INT);
@SuppressWarnings("rawtypes")
ForwardIndexReader denseFwd = denseFreeform.getForwardIndex();
@@ -184,7 +189,7 @@ public class OpenStructSparseDenseParityTest {
} else {
assertEquals(denseFwd.getInt(10, denseCtx), 42);
}
- assertEquals(sparseFwd.getString(10, sparseCtx), "42");
+ assertEquals(sparseFwd.getInt(10, sparseCtx), 42);
} finally {
dense.destroy();
sparse.destroy();
diff --git
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/V1Constants.java
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/V1Constants.java
index 5720bb4c125..c0fdf9a3717 100644
---
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/V1Constants.java
+++
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/V1Constants.java
@@ -213,6 +213,13 @@ public class V1Constants {
// built before multi-value keys — readers must treat an unlisted sparse
key as single-value.
public static final String SPARSE_MULTI_VALUE_KEYS =
"sparseMultiValueKeys";
+ // Value type of each key in this OPEN_STRUCT column's sparse column,
stored as a JSON object
+ // string mapping key name to DataType name. A dense key carries its
type in its own column
+ // metadata; a sparse key's lives here, because the blob it is stored in
is JSON and carries no
+ // type of its own. Absent on segments built before this was recorded --
readers must fall back
+ // to STRING for an unlisted sparse key, which is what every sparse key
used to read as.
+ public static final String SPARSE_KEY_TYPES = "sparseKeyTypes";
+
/// Partition function, all optional
public static final String PARTITION_FUNCTION = "partitionFunction";
public static final String NUM_PARTITIONS = "numPartitions";
diff --git
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/FieldIndexConfigsUtil.java
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/FieldIndexConfigsUtil.java
index 3f4136a2691..3bc61ae3fbb 100644
---
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/FieldIndexConfigsUtil.java
+++
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/FieldIndexConfigsUtil.java
@@ -58,16 +58,20 @@ public class FieldIndexConfigsUtil {
}
/// Builds a [FieldIndexConfigs] for a single column directly from one
[FieldConfig], without a
- /// `TableConfig` or `Schema`. The dictionary entry is derived from
`fieldConfig.getEncodingType()`
- /// (RAW => disabled, otherwise default-enabled); every other index type is
read from the modern
- /// `fieldConfig.getIndexes()` JSON (keyed by index pretty name), falling
back to each type's default config.
- /// A `null` fieldConfig yields built-in defaults (dictionary enabled).
+ /// `TableConfig` or `Schema`. The dictionary and forward entries are both
derived from
+ /// `fieldConfig.getEncodingType()` (RAW => dictionary disabled and a raw
forward index, otherwise
+ /// dictionary default-enabled and a dict-encoded forward index); every
other index type is read from the
+ /// modern `fieldConfig.getIndexes()` JSON (keyed by index pretty name),
falling back to each type's default
+ /// config. A `null` fieldConfig yields built-in defaults (dictionary
enabled).
///
/// This reads only the modern `indexes` format and never legacy
`IndexingConfig` lists, so it suits
/// synthetic columns (e.g. OPEN_STRUCT materialized children) that exist in
no schema.
public static FieldIndexConfigs fromFieldConfig(@Nullable FieldConfig
fieldConfig, FieldSpec fieldSpec) {
FieldIndexConfigs.Builder builder = new FieldIndexConfigs.Builder();
- boolean rawEncoded = fieldConfig != null && fieldConfig.getEncodingType()
== FieldConfig.EncodingType.RAW;
+ FieldConfig.EncodingType encodingType =
+ fieldConfig != null && fieldConfig.getEncodingType() != null ?
fieldConfig.getEncodingType()
+ : FieldConfig.EncodingType.DICTIONARY;
+ boolean rawEncoded = encodingType == FieldConfig.EncodingType.RAW;
builder.add(StandardIndexes.dictionary(),
rawEncoded ? DictionaryIndexConfig.DISABLED :
DictionaryIndexConfig.DEFAULT);
JsonNode indexes = fieldConfig != null ? fieldConfig.getIndexes() : null;
@@ -75,11 +79,35 @@ public class FieldIndexConfigsUtil {
if (indexType.getId().equals(StandardIndexes.DICTIONARY_ID)) {
continue;
}
+ if (indexType.getId().equals(StandardIndexes.FORWARD_ID)) {
+ // The forward index carries an encoding of its own, and
[ForwardIndexConfig]'s JSON default is
+ // DICTIONARY whatever the column-level encoding says. Left alone, a
RAW column would ask for a
+ // dict-encoded forward index over the dictionary disabled above --
which a reload tries to honour and
+ // fails on ("Dictionary should still exist after rebuilding
dict-encoded forward index"). Keep the two
+ // in step.
+ builder.add(StandardIndexes.forward(), forwardConfig(indexType,
indexes, encodingType));
+ continue;
+ }
addConfigFromIndexes(builder, indexType, indexes);
}
return builder.build();
}
+ private static ForwardIndexConfig forwardConfig(IndexType<?, ?, ?>
forwardIndexType, @Nullable JsonNode indexes,
+ FieldConfig.EncodingType encodingType) {
+ JsonNode node = indexes != null ?
indexes.get(forwardIndexType.getPrettyName()) : null;
+ if (node == null) {
+ return ForwardIndexConfig.getDefault(encodingType);
+ }
+ ForwardIndexConfig configured;
+ try {
+ configured = JsonUtils.jsonNodeToObject(node, ForwardIndexConfig.class);
+ } catch (IOException e) {
+ throw new IllegalArgumentException("Failed to parse 'forward' index
config from FieldConfig", e);
+ }
+ return new ForwardIndexConfig.Builder(configured, encodingType).build();
+ }
+
private static <C extends IndexConfig> void
addConfigFromIndexes(FieldIndexConfigs.Builder builder,
IndexType<C, ?, ?> indexType, @Nullable JsonNode indexes) {
JsonNode node = indexes != null ? indexes.get(indexType.getPrettyName()) :
null;
diff --git
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/metadata/ColumnMetadataImpl.java
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/metadata/ColumnMetadataImpl.java
index 8141d6725b9..cf866795b1d 100644
---
a/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/metadata/ColumnMetadataImpl.java
+++
b/pinot-segment-spi/src/main/java/org/apache/pinot/segment/spi/index/metadata/ColumnMetadataImpl.java
@@ -112,6 +112,8 @@ public class ColumnMetadataImpl implements ColumnMetadata {
@Nullable
private final Map<String, Integer> _sparseMultiValueKeys;
@Nullable
+ private final Map<String, DataType> _sparseKeyTypes;
+ @Nullable
private final CompressionMetadata _compressionMetadata;
/// Packed index sizes: the high 16 bits identify the index type and the low
48 bits hold its size.
@@ -126,7 +128,8 @@ public class ColumnMetadataImpl implements ColumnMetadata {
int totalNumberOfEntries, int maxNumberOfMultiValues, int
maxRowLengthInBytes, int bitsPerElement,
@Nullable PartitionFunction partitionFunction, @Nullable Set<Integer>
partitions, boolean autoGenerated,
@Nullable String parentColumn, @Nullable List<String> sparseKeys,
- @Nullable Map<String, Integer> sparseMultiValueKeys, @Nullable
CompressionMetadata compressionMetadata) {
+ @Nullable Map<String, Integer> sparseMultiValueKeys, @Nullable
Map<String, DataType> sparseKeyTypes,
+ @Nullable CompressionMetadata compressionMetadata) {
_fieldSpec = fieldSpec;
_totalDocs = totalDocs;
_cardinality = cardinality;
@@ -150,6 +153,7 @@ public class ColumnMetadataImpl implements ColumnMetadata {
_parentColumn = parentColumn;
_sparseKeys = sparseKeys;
_sparseMultiValueKeys = sparseMultiValueKeys;
+ _sparseKeyTypes = sparseKeyTypes;
_compressionMetadata = compressionMetadata;
}
@@ -283,6 +287,15 @@ public class ColumnMetadataImpl implements ColumnMetadata {
return _sparseMultiValueKeys;
}
+ /// The value type the segment build resolved for each key of this
OPEN_STRUCT column's sparse blob, or null on a
+ /// segment built before types were recorded. Only set on OPEN_STRUCT parent
columns. Which tier a key lands on is
+ /// a tuning decision, so it must not decide the key's type either: a dense
key reads with the type in its own
+ /// column metadata, and without this a sparse key would read as STRING no
matter what it holds.
+ @Nullable
+ public Map<String, DataType> getSparseKeyTypes() {
+ return _sparseKeyTypes;
+ }
+
@Override
public long getIndexSizeFor(IndexType type) {
if (_indexTypeSizes == null) {
@@ -383,6 +396,7 @@ public class ColumnMetadataImpl implements ColumnMetadata {
&& Objects.equals(_parentColumn, that._parentColumn)
&& Objects.equals(_sparseKeys, that._sparseKeys)
&& Objects.equals(_sparseMultiValueKeys, that._sparseMultiValueKeys)
+ && Objects.equals(_sparseKeyTypes, that._sparseKeyTypes)
&& Objects.equals(_compressionMetadata, that._compressionMetadata)
&& Objects.equals(_indexTypeSizes, that._indexTypeSizes);
}
@@ -392,8 +406,8 @@ public class ColumnMetadataImpl implements ColumnMetadata {
return Objects.hash(_fieldSpec, _totalDocs, _cardinality, _hasDictionary,
_forwardIndexEncoding, _sorted, _nonNull,
_minValue, _maxValue, _minMaxValueInvalid, _lengthOfShortestElement,
_lengthOfLongestElement, _isAscii,
_totalNumberOfEntries, _maxNumberOfMultiValues, _maxRowLengthInBytes,
_bitsPerElement, _partitionFunction,
- _partitions, _autoGenerated, _parentColumn, _sparseKeys,
_sparseMultiValueKeys, _compressionMetadata,
- _indexTypeSizes);
+ _partitions, _autoGenerated, _parentColumn, _sparseKeys,
_sparseMultiValueKeys, _sparseKeyTypes,
+ _compressionMetadata, _indexTypeSizes);
}
/// Rejoins a JSON property that
[org.apache.commons.configuration2.convert.LegacyListDelimiterHandler]
@@ -429,6 +443,7 @@ public class ColumnMetadataImpl implements ColumnMetadata {
+ ", _parentColumn=" + _parentColumn
+ ", _sparseKeys=" + _sparseKeys
+ ", _sparseMultiValueKeys=" + _sparseMultiValueKeys
+ + ", _sparseKeyTypes=" + _sparseKeyTypes
+ ", _compressionMetadata=" + _compressionMetadata
+ ", _indexTypeSizes=" + _indexTypeSizes
+ '}';
@@ -483,6 +498,17 @@ public class ColumnMetadataImpl implements ColumnMetadata {
}
}
+ Object rawSparseKeyTypes = config.getProperty(Column.getKeyFor(column,
Column.SPARSE_KEY_TYPES));
+ if (rawSparseKeyTypes != null) {
+ String jsonStr = rejoinJson(rawSparseKeyTypes);
+ try {
+ builder.setSparseKeyTypes(
+ JsonUtils.stringToObject(jsonStr, new TypeReference<Map<String,
DataType>>() { }));
+ } catch (Exception e) {
+ throw new RuntimeException("Failed to parse sparse key type manifest:
" + jsonStr, e);
+ }
+ }
+
// Set min/max value
DataType storedType = fieldSpec.getDataType().getStoredType();
if (fieldSpec instanceof ComplexFieldSpec) {
@@ -781,6 +807,7 @@ public class ColumnMetadataImpl implements ColumnMetadata {
private String _parentColumn;
private List<String> _sparseKeys;
private Map<String, Integer> _sparseMultiValueKeys;
+ private Map<String, DataType> _sparseKeyTypes;
private long _uncompressedValueSizeInBytes = UNAVAILABLE;
private ChunkCompressionType _forwardIndexChunkCompressionType;
private long _dictionaryUncompressedValueSizeInBytes = UNAVAILABLE;
@@ -900,6 +927,11 @@ public class ColumnMetadataImpl implements ColumnMetadata {
return this;
}
+ public Builder setSparseKeyTypes(Map<String, DataType> sparseKeyTypes) {
+ _sparseKeyTypes = sparseKeyTypes;
+ return this;
+ }
+
public Builder setSparseMultiValueKeys(Map<String, Integer>
sparseMultiValueKeys) {
_sparseMultiValueKeys = sparseMultiValueKeys;
return this;
@@ -965,6 +997,7 @@ public class ColumnMetadataImpl implements ColumnMetadata {
_lengthOfLongestElement, _isAscii, _totalNumberOfEntries,
_maxNumberOfMultiValues, _maxRowLengthInBytes,
_bitsPerElement, _partitionFunction, _partitions, _autoGenerated,
_parentColumn, _sparseKeys,
_sparseMultiValueKeys,
+ _sparseKeyTypes,
CompressionMetadata.create(_uncompressedValueSizeInBytes,
_forwardIndexChunkCompressionType,
_dictionaryUncompressedValueSizeInBytes));
}
diff --git
a/pinot-segment-spi/src/test/java/org/apache/pinot/segment/spi/index/metadata/ColumnMetadataSparseKeysTest.java
b/pinot-segment-spi/src/test/java/org/apache/pinot/segment/spi/index/metadata/ColumnMetadataSparseKeysTest.java
index bf1dcc131a5..3195546a15d 100644
---
a/pinot-segment-spi/src/test/java/org/apache/pinot/segment/spi/index/metadata/ColumnMetadataSparseKeysTest.java
+++
b/pinot-segment-spi/src/test/java/org/apache/pinot/segment/spi/index/metadata/ColumnMetadataSparseKeysTest.java
@@ -20,6 +20,7 @@ package org.apache.pinot.segment.spi.index.metadata;
import java.io.File;
import java.util.List;
+import java.util.Map;
import org.apache.commons.configuration2.PropertiesConfiguration;
import org.apache.commons.configuration2.convert.LegacyListDelimiterHandler;
import org.apache.pinot.segment.spi.V1Constants.MetadataKeys.Column;
@@ -68,6 +69,34 @@ public class ColumnMetadataSparseKeysTest {
return CommonsConfigurationUtils.fromFile(_tempFile);
}
+ /// The type manifest is the sparse tier's answer to a dense child's own
column metadata, so it has to survive
+ /// the same save/load round trip the key manifest does -- the commas
between entries are exactly what
+ /// LegacyListDelimiterHandler fragments.
+ @Test
+ public void testSparseKeyTypesRoundTripThroughSaveLoad() throws Exception {
+ PropertiesConfiguration config = baseParentProps();
+ config.setProperty(Column.getKeyFor(COLUMN, Column.SPARSE_KEY_TYPES),
+ "{\"region\":\"STRING\",\"latencyMs\":\"LONG\",\"freeform\":\"INT\"}");
+
+ PropertiesConfiguration reloaded = saveAndReload(config);
+ ColumnMetadataImpl metadata =
ColumnMetadataImpl.fromPropertiesConfiguration(reloaded, 10, COLUMN);
+ assertEquals(metadata.getSparseKeyTypes(),
+ Map.of("region", DataType.STRING, "latencyMs", DataType.LONG,
"freeform", DataType.INT));
+ }
+
+ /// A segment built before types were recorded carries no manifest, and must
load rather than fail -- readers
+ /// fall back to the STRING every sparse key used to read as.
+ @Test
+ public void testSparseKeyTypesAbsentOnOlderSegment() throws Exception {
+ PropertiesConfiguration config = baseParentProps();
+ config.setProperty(Column.getKeyFor(COLUMN, Column.SPARSE_KEYS),
"[\"region\",\"latencyMs\"]");
+
+ PropertiesConfiguration reloaded = saveAndReload(config);
+ ColumnMetadataImpl metadata =
ColumnMetadataImpl.fromPropertiesConfiguration(reloaded, 10, COLUMN);
+ assertEquals(metadata.getSparseKeys(), List.of("region", "latencyMs"));
+ assertNull(metadata.getSparseKeyTypes());
+ }
+
@Test
public void testSparseKeysListRoundTripsThroughSaveLoad() throws Exception {
PropertiesConfiguration config = baseParentProps();
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]