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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -167,16 +167,14 @@ private static void updateSegmentZKMetadata(String tableNameWithType, SegmentZKM

// Set partition metadata
Map<String, ColumnPartitionMetadata> columnPartitionMap = new HashMap<>();
for (Map.Entry<String, ColumnMetadata> entry : segmentMetadata.getColumnMetadataMap().entrySet()) {
ColumnMetadata columnMetadata = entry.getValue();
segmentMetadata.forEachColumn((column, columnMetadata) -> {
PartitionFunction partitionFunction = columnMetadata.getPartitionFunction();
if (partitionFunction != null) {
ColumnPartitionMetadata columnPartitionMetadata =
columnPartitionMap.put(column,
new ColumnPartitionMetadata(partitionFunction.getName(), partitionFunction.getNumPartitions(),
columnMetadata.getPartitions(), partitionFunction.getFunctionConfig());
columnPartitionMap.put(entry.getKey(), columnPartitionMetadata);
columnMetadata.getPartitions(), partitionFunction.getFunctionConfig()));
}
}
});
segmentZKMetadata.setPartitionMetadata(
!columnPartitionMap.isEmpty() ? new SegmentPartitionMetadata(columnPartitionMap) : null);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import java.util.TreeMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.BiConsumer;
import org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
import org.apache.pinot.common.partition.function.MurmurPartitionFunction;
import org.apache.pinot.segment.spi.ColumnMetadata;
Expand All @@ -30,6 +31,9 @@
import org.joda.time.Interval;
import org.mockito.Mockito;

import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;

Expand All @@ -40,6 +44,22 @@ public class SegmentMetadataMockUtils {
private SegmentMetadataMockUtils() {
}

/// Stubs the column accessors of a mocked segment metadata. Mockito stubs every method, so stubbing only
/// `getColumnMetadataMap()` leaves the accessors production code reads (the real implementation holds sorted
/// arrays, not a map) answering `null`.
private static void stubColumns(SegmentMetadata segmentMetadata, TreeMap<String, ColumnMetadata> columns) {
when(segmentMetadata.getColumnMetadataMap()).thenReturn(columns);
when(segmentMetadata.getAllColumns()).thenReturn(columns.navigableKeySet());
when(segmentMetadata.getAllColumnMetadata()).thenReturn(columns.values());
when(segmentMetadata.getNumColumns()).thenReturn(columns.size());
when(segmentMetadata.getColumnMetadataFor(anyString())).thenAnswer(
call -> columns.get(call.<String>getArgument(0)));
doAnswer(call -> {
columns.forEach(call.<BiConsumer<String, ColumnMetadata>>getArgument(0));
return null;
}).when(segmentMetadata).forEachColumn(any());
}

public static SegmentMetadata mockSegmentMetadata(String tableName, String segmentName, int numTotalDocs,
String crc, long startTime, long endTime, TimeUnit timeUnit) {
SegmentMetadata segmentMetadata = Mockito.mock(SegmentMetadata.class);
Expand Down Expand Up @@ -100,7 +120,7 @@ public static SegmentMetadata mockSegmentMetadata(String tableName, String segme
when(colMeta.getPartitionFunction()).thenReturn(new MurmurPartitionFunction(numPartitions, null));
TreeMap<String, ColumnMetadata> columnMetadataMap = new TreeMap<>();
columnMetadataMap.put(partitionColumn, colMeta);
when(segmentMetadata.getColumnMetadataMap()).thenReturn(columnMetadataMap);
stubColumns(segmentMetadata, columnMetadataMap);
return segmentMetadata;
}

Expand All @@ -112,17 +132,14 @@ public static SegmentMetadata mockSegmentMetadataWithPartitionInfo(String rawTab
when(columnMetadata.getPartitionFunction()).thenReturn(new MurmurPartitionFunction(5, null));

SegmentMetadataImpl segmentMetadata = mock(SegmentMetadataImpl.class);
if (columnName != null) {
when(segmentMetadata.getColumnMetadataFor(columnName)).thenReturn(columnMetadata);
}
when(segmentMetadata.getTableName()).thenReturn(rawTableName);
when(segmentMetadata.getName()).thenReturn(segmentName);
when(segmentMetadata.getCrc()).thenReturn("0");
when(segmentMetadata.getDataCrc()).thenReturn("1");

TreeMap<String, ColumnMetadata> columnMetadataMap = new TreeMap<>();
columnMetadataMap.put(columnName, columnMetadata);
when(segmentMetadata.getColumnMetadataMap()).thenReturn(columnMetadataMap);
stubColumns(segmentMetadata, columnMetadataMap);
return segmentMetadata;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -312,7 +312,7 @@ protected double getMaxDedupTime(IndexSegment segment) {
// so far to process this segment, as mutable segment is always considered to be within TTL
return _largestSeenTime.get();
}
return ((Number) segment.getSegmentMetadata().getColumnMetadataMap().get(_dedupTimeColumn)
return ((Number) segment.getSegmentMetadata().getColumnMetadataFor(_dedupTimeColumn)
.getMaxValue()).doubleValue();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@
package org.apache.pinot.segment.local.indexsegment.immutable;

import com.google.common.base.Preconditions;
import java.util.Collections;
import java.util.List;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
Expand Down Expand Up @@ -80,16 +79,16 @@ public SegmentMetadataImpl getSegmentMetadata() {
return _segmentMetadata;
}

// Both are views of the column metadata map, so neither builds the segment schema (see SegmentMetadataImpl)
// Both are views of the segment's column metadata, so neither builds the segment schema (see SegmentMetadataImpl)

@Override
public Set<String> getColumnNames() {
return Collections.unmodifiableSet(_segmentMetadata.getColumnMetadataMap().keySet());
return _segmentMetadata.getAllColumns();
}

@Override
public Set<String> getPhysicalColumnNames() {
return new PhysicalColumnNames(_segmentMetadata.getColumnMetadataMap());
return new PhysicalColumnNames(_segmentMetadata);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,13 +25,11 @@
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.TreeMap;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicBoolean;
Expand Down Expand Up @@ -156,19 +154,15 @@ public ImmutableSegmentImpl(
_columnMaterializer = null;
_openStructChildren = null;
_materializationLock = null;
TreeMap<String, ColumnMetadata> columnMetadataMap = segmentMetadata.getColumnMetadataMap();
_columnNames = Collections.unmodifiableSet(columnMetadataMap.keySet());
_physicalColumnNames = new PhysicalColumnNames(columnMetadataMap);
_dataSources = new Object2ObjectOpenHashMap<>(columnMetadataMap.size());
_columnNames = segmentMetadata.getAllColumns();
_physicalColumnNames = new PhysicalColumnNames(segmentMetadata);
_dataSources = new Object2ObjectOpenHashMap<>(segmentMetadata.getNumColumns());

Map<String, Map<String, DataSource>> openStructDenseChildren = new HashMap<>();
Map<String, DataSource> openStructSparseChildren = new HashMap<>();
Set<String> openStructParents = new HashSet<>();

for (Map.Entry<String, ColumnMetadata> entry : columnMetadataMap.entrySet()) {
String colName = entry.getKey();
ColumnMetadata columnMetadata = entry.getValue();

segmentMetadata.forEachColumn((colName, columnMetadata) -> {
if (columnMetadata instanceof ColumnMetadataImpl && ((ColumnMetadataImpl) columnMetadata).isMaterializedChild()) {
String parent = ((ColumnMetadataImpl) columnMetadata).getParentColumn();
openStructParents.add(parent);
Expand All @@ -179,19 +173,19 @@ public ImmutableSegmentImpl(
openStructDenseChildren.computeIfAbsent(parent, k -> new HashMap<>())
.put(OpenStructNaming.parseKey(colName), childDs);
}
continue;
return;
}

if (columnMetadata.getFieldSpec().getDataType() == FieldSpec.DataType.MAP) {
_dataSources.put(colName, new ImmutableMapDataSource(columnMetadata, _indexContainerMap.get(colName)));
} else {
_dataSources.put(colName, new ImmutableDataSource(columnMetadata, _indexContainerMap.get(colName)));
}
}
});

for (String parent : openStructParents) {
// The parent's spec comes from its column metadata, not from the segment schema (see _columnNames)
ColumnMetadata parentMetadata = columnMetadataMap.get(parent);
ColumnMetadata parentMetadata = segmentMetadata.getColumnMetadataFor(parent);
FieldSpec fieldSpec = parentMetadata != null ? parentMetadata.getFieldSpec() : null;
if (!(fieldSpec instanceof ComplexFieldSpec)) {
continue;
Expand Down Expand Up @@ -231,9 +225,8 @@ public ImmutableSegmentImpl(
_columnMaterializer = columnMaterializer;
_openStructChildren = groupOpenStructChildren(segmentMetadata);
_materializationLock = new ReentrantReadWriteLock();
TreeMap<String, ColumnMetadata> columnMetadataMap = segmentMetadata.getColumnMetadataMap();
_columnNames = Collections.unmodifiableSet(columnMetadataMap.keySet());
_physicalColumnNames = new PhysicalColumnNames(columnMetadataMap);
_columnNames = segmentMetadata.getAllColumns();
_physicalColumnNames = new PhysicalColumnNames(segmentMetadata);
_dataSources = new ConcurrentHashMap<>();
for (String column : materializedIndexContainers.keySet()) {
materializeDataSource(column);
Expand All @@ -244,21 +237,14 @@ public ImmutableSegmentImpl(
/// metadata declares them complex (the same rule the eager constructor applies).
@Nullable
private static Map<String, List<String>> groupOpenStructChildren(SegmentMetadataImpl segmentMetadata) {
Map<String, List<String>> children = null;
Map<String, ColumnMetadata> columnMetadataMap = segmentMetadata.getColumnMetadataMap();
for (Map.Entry<String, ColumnMetadata> entry : columnMetadataMap.entrySet()) {
if (entry.getValue() instanceof ColumnMetadataImpl impl && impl.isMaterializedChild()) {
if (children == null) {
children = new HashMap<>();
}
children.computeIfAbsent(impl.getParentColumn(), k -> new ArrayList<>()).add(entry.getKey());
Map<String, List<String>> children = new HashMap<>();
segmentMetadata.forEachColumn((column, columnMetadata) -> {
if (columnMetadata instanceof ColumnMetadataImpl impl && impl.isMaterializedChild()) {
children.computeIfAbsent(impl.getParentColumn(), k -> new ArrayList<>()).add(column);
}
}
if (children == null) {
return null;
}
});
children.keySet().removeIf(parent -> {
ColumnMetadata parentMetadata = columnMetadataMap.get(parent);
ColumnMetadata parentMetadata = segmentMetadata.getColumnMetadataFor(parent);
return parentMetadata == null || !(parentMetadata.getFieldSpec() instanceof ComplexFieldSpec);
});
return children.isEmpty() ? null : children;
Expand All @@ -268,7 +254,7 @@ private static Map<String, List<String>> groupOpenStructChildren(SegmentMetadata
/// such column. OPEN_STRUCT child columns are reachable only through their parent, as in the eager mode.
@Nullable
private DataSource materializeDataSource(String column) {
ColumnMetadata columnMetadata = _segmentMetadata.getColumnMetadataMap().get(column);
ColumnMetadata columnMetadata = _segmentMetadata.getColumnMetadataFor(column);
boolean openStructParent = _openStructChildren != null && _openStructChildren.containsKey(column);
if (!openStructParent && (columnMetadata == null || isMaterializedChild(columnMetadata))) {
return null;
Expand All @@ -294,11 +280,10 @@ private DataSource createDataSource(String column, ColumnMetadata columnMetadata
}

private DataSource createOpenStructDataSource(String parent) {
Map<String, ColumnMetadata> columnMetadataMap = _segmentMetadata.getColumnMetadataMap();
Map<String, DataSource> denseChildren = new HashMap<>();
DataSource sparseChild = null;
for (String child : _openStructChildren.get(parent)) {
ColumnMetadata childMetadata = columnMetadataMap.get(child);
ColumnMetadata childMetadata = _segmentMetadata.getColumnMetadataFor(child);
DataSource childDataSource =
new ImmutableDataSource(childMetadata, materializedIndexContainer(child, childMetadata));
if (OpenStructNaming.isSparseColumn(child)) {
Expand All @@ -307,7 +292,7 @@ private DataSource createOpenStructDataSource(String parent) {
denseChildren.put(OpenStructNaming.parseKey(child), childDataSource);
}
}
ColumnMetadata parentMetadata = columnMetadataMap.get(parent);
ColumnMetadata parentMetadata = _segmentMetadata.getColumnMetadataFor(parent);
ComplexFieldSpec fieldSpec = (ComplexFieldSpec) parentMetadata.getFieldSpec();
List<String> sparseKeys = parentMetadata instanceof ColumnMetadataImpl impl ? impl.getSparseKeys() : null;
return new ImmutableOpenStructDataSource(fieldSpec, denseChildren, sparseChild, _segmentMetadata.getTotalDocs(),
Expand Down Expand Up @@ -412,7 +397,7 @@ public boolean isReloadNeeded(IndexLoadingConfig indexLoadingConfig)
public <I extends IndexReader> I getIndex(String column, IndexType<?, I, ?> type) {
ColumnIndexContainer container = _indexContainerMap.get(column);
if (container == null && _columnMaterializer != null) {
ColumnMetadata columnMetadata = _segmentMetadata.getColumnMetadataMap().get(column);
ColumnMetadata columnMetadata = _segmentMetadata.getColumnMetadataFor(column);
if (columnMetadata != null) {
Lock lock = _materializationLock.readLock();
lock.lock();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,6 @@
import org.apache.pinot.segment.local.segment.virtualcolumn.VirtualColumnProviderFactory;
import org.apache.pinot.segment.local.startree.v2.store.StarTreeIndexContainer;
import org.apache.pinot.segment.local.utils.SegmentOperationsThrottlerSet;
import org.apache.pinot.segment.spi.ColumnMetadata;
import org.apache.pinot.segment.spi.ImmutableSegment;
import org.apache.pinot.segment.spi.converter.SegmentFormatConverter;
import org.apache.pinot.segment.spi.creator.SegmentVersion;
Expand Down Expand Up @@ -211,9 +210,8 @@ public static ImmutableSegment load(SegmentDirectory segmentDirectory, IndexLoad
}

// Remove columns not in schema from the metadata
Map<String, ColumnMetadata> columnMetadataMap = segmentMetadata.getColumnMetadataMap();
if (schema != null) {
Set<String> columnsInMetadata = new HashSet<>(columnMetadataMap.keySet());
Set<String> columnsInMetadata = new HashSet<>(segmentMetadata.getAllColumns());
columnsInMetadata.removeIf(schema::hasColumn);
// Materialized OPEN_STRUCT child columns (col$key, col$__sparse__) live in segment metadata
// but not in the user-facing schema. Keep them when the parent OPEN_STRUCT column is in the
Expand All @@ -233,7 +231,7 @@ public static ImmutableSegment load(SegmentDirectory segmentDirectory, IndexLoad
}
}
} else {
indexLoadingConfig.addKnownColumns(columnMetadataMap.keySet());
indexLoadingConfig.addKnownColumns(segmentMetadata.getAllColumns());
}

SegmentDirectory.Reader segmentReader = segmentDirectory.createReader();
Expand All @@ -245,11 +243,13 @@ public static ImmutableSegment load(SegmentDirectory segmentDirectory, IndexLoad
return segment;
}

Map<String, ColumnIndexContainer> indexContainerMap = new Object2ObjectOpenHashMap<>(columnMetadataMap.size());
for (Map.Entry<String, ColumnMetadata> entry : columnMetadataMap.entrySet()) {
Map<String, ColumnIndexContainer> indexContainerMap =
new Object2ObjectOpenHashMap<>(segmentMetadata.getNumColumns());
for (String column : segmentMetadata.getAllColumns()) {
// FIXME: text-index only works with local SegmentDirectory
indexContainerMap.put(entry.getKey(),
new PhysicalColumnIndexContainer(segmentReader, entry.getValue(), indexLoadingConfig));
indexContainerMap.put(column,
new PhysicalColumnIndexContainer(segmentReader, segmentMetadata.getColumnMetadataFor(column),
indexLoadingConfig));
}

instantiateVirtualColumns(segmentMetadata, indexContainerMap);
Expand Down Expand Up @@ -287,14 +287,13 @@ private static ImmutableSegmentImpl loadWithLazyColumns(SegmentDirectory segment
SegmentDirectory.Reader segmentReader, SegmentMetadataImpl segmentMetadata,
IndexLoadingConfig indexLoadingConfig)
throws IOException {
Map<String, ColumnMetadata> columnMetadataMap = segmentMetadata.getColumnMetadataMap();
MultiColumnLuceneTextIndexReader mcTextReader = null;
Set<String> mcTextColumns = Set.of();
if (segmentReader.hasMultiColumnTextIndex()) {
mcTextReader = new MultiColumnLuceneTextIndexReader(segmentMetadata);
mcTextColumns = Set.copyOf(segmentMetadata.getMultiColumnTextMetadata().getColumns());
}
ColumnMaterializer columnMaterializer = new ColumnMaterializer(segmentReader, columnMetadataMap.keySet(),
ColumnMaterializer columnMaterializer = new ColumnMaterializer(segmentReader, segmentMetadata.getAllColumns(),
indexLoadingConfig.getFieldIndexConfigByColName(), indexLoadingConfig.isForwardIndexOnly(), mcTextReader,
mcTextColumns);

Expand All @@ -305,7 +304,7 @@ private static ImmutableSegmentImpl loadWithLazyColumns(SegmentDirectory segment
if (segmentReader.hasStarTreeIndex()) {
starTreeIndexContainer = new StarTreeIndexContainer(segmentReader, segmentMetadata,
column -> indexContainerMap.computeIfAbsent(column,
k -> columnMaterializer.createIndexContainer(columnMetadataMap.get(k))));
k -> columnMaterializer.createIndexContainer(segmentMetadata.getColumnMetadataFor(k))));
}

return new ImmutableSegmentImpl(segmentDirectory, segmentMetadata, columnMaterializer, indexContainerMap,
Expand All @@ -314,24 +313,23 @@ private static ImmutableSegmentImpl loadWithLazyColumns(SegmentDirectory segment

/// Creates the index containers and column metadata of the built-in virtual columns and registers them in the
/// segment metadata. Registering the metadata is what makes the segment schema include the virtual columns: the
/// schema is derived from the column metadata map on demand ([SegmentMetadataImpl#getSchema()]) and is deliberately
/// schema is derived from the column metadata on demand ([SegmentMetadataImpl#getSchema()]) and is deliberately
/// not built here, so a loaded segment retains no per-column schema entries until something asks for its schema.
/// A physical column of the same name wins, as in the schema-based registration this replaces.
private static void instantiateVirtualColumns(SegmentMetadataImpl segmentMetadata,
Map<String, ColumnIndexContainer> indexContainerMap) {
Map<String, ColumnMetadata> columnMetadataMap = segmentMetadata.getColumnMetadataMap();
String segmentName = segmentMetadata.getName();
for (BuiltInVirtualColumnDefinitions.Definition definition : BuiltInVirtualColumnDefinitions.DEFINITIONS) {
String columnName = definition.getName();
if (columnMetadataMap.containsKey(columnName)) {
if (segmentMetadata.getColumnMetadataFor(columnName) != null) {
continue;
}
FieldSpec fieldSpec = VirtualColumnProviderFactory.createBuiltInFieldSpec(definition, segmentName);
VirtualColumnContext context =
new VirtualColumnContext(fieldSpec, segmentMetadata.getTotalDocs(), segmentMetadata);
VirtualColumnProvider provider = VirtualColumnProviderFactory.buildProvider(context);
indexContainerMap.put(columnName, provider.buildColumnIndexContainer(context));
columnMetadataMap.put(columnName, provider.buildMetadata(context));
segmentMetadata.addColumnMetadata(columnName, provider.buildMetadata(context));
}
}

Expand Down
Loading
Loading