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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 41 additions & 0 deletions docs/docs/multimodal-table/global-index/btree.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,47 @@ print(pa_table)

</Tabs>

## Composite BTree Indexes

Spark and Flink can create a composite index by listing columns in key order:

```sql
CALL sys.create_global_index(
`table` => 'db.my_table',
index_column => 'category,item_number',
index_type => 'btree'
);

SELECT * FROM my_table
WHERE category = 'category-a'
AND item_number = 107;
```

The index stores typed tuples and one row-ID posting list for each distinct tuple.
Equality conditions on every key column use a composite point lookup, regardless
of their order in the SQL predicate. This avoids expanding each column's posting
list before intersecting them. Additional predicates remain row filters, and OR
branches can each select a composite lookup. When several definitions match,
selection prefers the one with more key columns.

Eligible composite point queries honor `global-index.query-in-reader.enabled`.
With reader evaluation enabled, supported index conditions narrow the candidate
rows and remaining conditions are evaluated against the data.

Composite indexes can coexist with single-column indexes on the same columns.
Composite lookup currently requires equality on all key columns; conditions on
only some key columns, ranges and IN predicates use available single-column
indexes or ordinary scans. Vector and full-text search pre-filters use
single-column indexes or their data fallback.

In `full` and `detail` modes, rows outside the selected composite index's coverage
remain eligible for data filtering. Calling the procedure again with the same
ordered column list fills missing index ranges. Under the `IGNORE` column-update
policy, updates to any key column invalidate the affected composite index ranges;
calling the procedure again rebuilds those ranges.
To drop the composite definition, pass the same ordered column list to
`sys.drop_global_index`; independent single-column definitions remain in place.

## BTree Options

| Option | Default | Description |
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.InternalRowUtils;
import org.apache.paimon.utils.Range;

import org.apache.flink.api.common.functions.Partitioner;
Expand Down Expand Up @@ -150,12 +151,34 @@ public static Optional<DataStream<Committable>> buildIndexStream(
PartitionPredicate partitionPredicate,
Options userOptions)
throws Exception {
List<List<String>> definitions = new ArrayList<>();
for (String name : indexColumns) {
definitions.add(Collections.singletonList(name));
}
return buildIndexDefinitions(
env,
indexScannerSupplier,
table,
definitions,
indexType,
partitionPredicate,
userOptions);
}

private static Optional<DataStream<Committable>> buildIndexDefinitions(
StreamExecutionEnvironment env,
Supplier<SortedGlobalIndexScanner> indexScannerSupplier,
FileStoreTable table,
List<List<String>> definitions,
String indexType,
PartitionPredicate partitionPredicate,
Options userOptions)
throws Exception {
List<DataStream<Committable>> allStreams = new ArrayList<>();
for (String indexColumn : indexColumns) {
for (List<String> indexColumns : definitions) {
String indexColumn = indexColumns.get(0);
SortedGlobalIndexScanner indexScanner =
indexScannerSupplier
.get()
.withIndexFields(Collections.singletonList(indexColumn));
indexScannerSupplier.get().withIndexFields(indexColumns);
if (partitionPredicate != null) {
indexScanner = indexScanner.withPartitionPredicate(partitionPredicate);
}
Expand All @@ -176,18 +199,12 @@ public static Optional<DataStream<Committable>> buildIndexStream(
continue;
}

// 2. Select necessary columns (index field + ROW_ID)
List<String> selectedColumns = new ArrayList<>();
selectedColumns.add(indexColumn);

RowType dataReadType =
SpecialFields.rowTypeWithRowId(table.rowType().project(selectedColumns));
// 2. Select the ordered index fields and ROW_ID.
RowType sourceReadType = table.rowType().project(indexColumns);
RowType dataReadType = SpecialFields.rowTypeWithRowId(sourceReadType);
String buildTaskIdField = buildTaskIdFieldName(dataReadType);
GlobalIndexer indexer =
GlobalIndexer.create(
indexType,
Collections.singletonList(table.rowType().getField(indexColumn)),
userOptions);
GlobalIndexer.create(indexType, sourceReadType.getFields(), userOptions);
if (!(indexer instanceof SortedGlobalIndexer)) {
throw new IllegalArgumentException(
"Index algorithm " + indexType + " does not expose sorted index keys.");
Expand All @@ -200,7 +217,6 @@ public static Optional<DataStream<Committable>> buildIndexStream(
int indexFieldPos = sortReadType.getFieldIndex(indexColumn);
int rowIdPos = sortReadType.getFieldIndex(SpecialFields.ROW_ID.name());
DataType indexFieldType = sortReadType.getTypeAt(indexFieldPos);
DataType sourceFieldType = table.rowType().getField(indexColumn).type();

// 3. Calculate maximum parallelism bound
long recordsPerRange =
Expand All @@ -220,7 +236,7 @@ public static Optional<DataStream<Committable>> buildIndexStream(
BinaryRow partition = partitionEntry.getKey();
byte[] partitionBytes = binaryRowSerializer.serializeToBytes(partition);
Map<Range, List<Split>> ranges =
keyExtractor.isIdentity()
keyExtractor.isIdentity() && sourceReadType.getFieldCount() == 1
? partitionEntry.getValue()
: shardSplitsByRowRange(partitionEntry.getValue(), recordsPerRange);
for (Map.Entry<Range, List<Split>> entry : ranges.entrySet()) {
Expand Down Expand Up @@ -249,14 +265,14 @@ public static Optional<DataStream<Committable>> buildIndexStream(
splitTasks,
readBuilder,
new SortedGlobalIndexWriter(table, indexType, userOptions)
.withIndexFields(Collections.singletonList(indexColumn)),
.withIndexFields(indexColumns),
scanResult.scanSnapshotId(),
partitionFieldSize,
taskIdPos,
indexFieldPos,
rowIdPos,
indexFieldType,
sourceFieldType,
sourceReadType,
keyExtractor,
coreOptions,
sortReadType,
Expand Down Expand Up @@ -306,19 +322,22 @@ private static List<Committable> createDeleteCommittables(
public static void buildIndexAndExecute(
StreamExecutionEnvironment env,
FileStoreTable table,
String indexColumn,
List<String> indexColumns,
String indexType,
PartitionPredicate partitionPredicate,
Options userOptions)
throws Exception {
if (buildIndex(
env,
() -> new SortedGlobalIndexScanner(table, indexType, userOptions),
table,
Collections.singletonList(indexColumn),
indexType,
partitionPredicate,
userOptions)) {
Optional<DataStream<Committable>> written =
buildIndexDefinitions(
env,
() -> new SortedGlobalIndexScanner(table, indexType, userOptions),
table,
Collections.singletonList(indexColumns),
indexType,
partitionPredicate,
userOptions);
if (written.isPresent()) {
commit(table, written.get(), CoreOptions.createCommitUser(userOptions));
env.execute("Create " + indexType + " global index for table: " + table.name());
}
}
Expand All @@ -335,7 +354,7 @@ protected static DataStream<Committable> executeForBuildTasks(
int indexFieldPos,
int rowIdPos,
DataType indexFieldType,
DataType sourceFieldType,
RowType sourceReadType,
GlobalIndexKeyExtractor keyExtractor,
CoreOptions coreOptions,
RowType readType,
Expand All @@ -354,14 +373,14 @@ protected static DataStream<Committable> executeForBuildTasks(
.transform(
"Read Data",
InternalTypeInfo.fromRowType(readType),
new ReadDataOperator(readBuilder, keyExtractor, sourceFieldType))
new ReadDataOperator(readBuilder, keyExtractor, sourceReadType))
.setParallelism(parallelism);

DataStream<InternalRow> sortedStream =
sortRows(
env,
rowDataStream,
keyExtractor.isIdentity(),
keyExtractor.isIdentity() && sourceReadType.getFieldCount() == 1,
taskIdPos,
indexFieldPos,
coreOptions,
Expand Down Expand Up @@ -516,25 +535,34 @@ private static class ReadDataOperator

private final ReadBuilder readBuilder;
private final GlobalIndexKeyExtractor keyExtractor;
private final DataType sourceFieldType;
private final RowType sourceReadType;

private transient TableRead tableRead;
private transient InternalRow.FieldGetter sourceFieldGetter;
private transient InternalRow.FieldGetter[] compositeGetters;

public ReadDataOperator(
ReadBuilder readBuilder,
GlobalIndexKeyExtractor keyExtractor,
DataType sourceFieldType) {
RowType sourceReadType) {
this.readBuilder = readBuilder;
this.keyExtractor = keyExtractor;
this.sourceFieldType = sourceFieldType;
this.sourceReadType = sourceReadType;
}

@Override
public void open() throws Exception {
super.open();
this.tableRead = readBuilder.newRead();
this.sourceFieldGetter = InternalRow.createFieldGetter(sourceFieldType, 0);
if (sourceReadType.getFieldCount() > 1) {
compositeGetters = new InternalRow.FieldGetter[sourceReadType.getFieldCount()];
for (int i = 0; i < compositeGetters.length; i++) {
compositeGetters[i] =
InternalRow.createFieldGetter(sourceReadType.getTypeAt(i), i);
}
} else {
sourceFieldGetter = InternalRow.createFieldGetter(sourceReadType.getTypeAt(0), 0);
}
}

@Override
Expand All @@ -546,10 +574,20 @@ public void processElement(StreamRecord<SortedSplitTask> element) throws Excepti
try {
InternalRow row;
while ((row = batch.next()) != null) {
long rowId = row.getLong(1);
long rowId = row.getLong(sourceReadType.getFieldCount());
Object sourceValue;
if (compositeGetters == null) {
sourceValue = sourceFieldGetter.getFieldOrNull(row);
} else {
GenericRow tuple = new GenericRow(compositeGetters.length);
for (int i = 0; i < compositeGetters.length; i++) {
tuple.setField(i, compositeGetters[i].getFieldOrNull(row));
}
sourceValue = InternalRowUtils.copy(tuple, sourceReadType);
}
boolean[] emitted = new boolean[1];
keyExtractor.extract(
sourceFieldGetter.getFieldOrNull(row),
sourceValue,
key -> {
emitted[0] = true;
output.collect(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -135,7 +135,7 @@ public String[] call(
SortedIndexTopoBuilder.buildIndexAndExecute(
procedureContext.getExecutionEnvironment(),
table,
indexColumns.get(0),
indexColumns,
indexType,
partitionPredicate,
userOptions);
Expand Down
Loading
Loading