Skip to content
Closed
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
76 changes: 75 additions & 1 deletion docs/docs/multimodal-table/global-index/btree.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,80 @@ print(pa_table)

</Tabs>

## Composite BTree Indexes

Spark and Flink can build 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 in lexicographic column order and one row-ID posting
list for each distinct tuple. Conditions are matched by column name, regardless of
the SQL condition order. For an index on `(category, item_number, tag)`:

| Conditions | Index access |
|---|---|
| Equality on all key columns | Point lookup; avoids expanding individual columns' posting lists. |
| `category = 'category-a'` | Scan the matching leading prefix. |
| `category = 'category-a' AND item_number > 107` | Equality prefix followed by a range. `>=`, `<`, `<=`, and `BETWEEN` are also supported. |
| `category IN ('category-a', 'category-b') AND item_number BETWEEN 10 AND 20` | One range for each equality/IN prefix. |
| `category IS NULL AND item_number = 107` | NULL is a fixed prefix value. `IS NOT NULL` can bound a range. |
| `category = 'category-a' AND item_number > 107 AND tag = 'x'` | Scan the prefix/range and test `tag` on index keys before decoding row IDs. |
| `category = 'category-a' AND tag = 'x'` | Scan the category prefix; test `tag` on index keys. The gap prevents `tag` from narrowing the scan interval. |
| Conditions on only `item_number` or `tag` | Use available single-column indexes or ordinary scans. No skip scan is performed. |

Leading `=`, `IN`, and `IS NULL` conditions extend the lookup prefix. The first
range column ends interval construction. Supported conditions on later index
columns are still evaluated on keys before expanding postings; they do not
narrow the scanned interval. Other conditions remain data filters. OR branches
can each choose their own index path. String predicates such as `LIKE` can filter
keys within a selected prefix/range; they do not construct composite prefix
ranges by themselves.

IN expansion is limited to 256 intervals per composite index lookup. The limit applies
to the product of distinct, matching values in consecutive key columns.
Exceeding the limit falls back to available single-column indexes or ordinary
scans. Prefix and range scans also obey `btree-index.fallback-scan-max-size`,
using the total file bytes remaining after min/max pruning, per index row-range
group. Complete point lookups do not use this scan budget. Budget-ineligible
composite paths are excluded before selecting an index, preserving available
single-column alternatives.

Selection uses metadata heuristics: more bounded leading columns, more indexed
filter columns, broader row-ID coverage, complete point probes, fewer intervals,
then fewer selected file bytes and fewer total key columns. A single-column
BTree is preferred when only the composite leading column is filtered and the
single-column index covers all of the composite candidate's row ranges. This is
not a statistics-based cost optimizer.

Composite prefix/range queries and queries combining composite and scalar
conditions, including scalar alternatives after a composite path is rejected,
are evaluated during planning, so runtime support and scan budgets
are resolved before data-split pruning. Pure composite point queries, including
bounded IN combinations and NULL prefixes, can use reader-side evaluation.

Vector and full-text search pre-filters currently use single-column indexes or
the usual data fallback; they do not evaluate composite BTree indexes.
Composite indexes can coexist with single-column indexes on the same columns.

Coverage is determined by the selected index. In `full` and `detail` modes, rows
outside the composite index's coverage are scanned even when single-column
indexes cover those rows. Incremental builds fill missing ranges; under the
`IGNORE` column-update policy, updates to any key column refresh the affected
composite index ranges. Dropping an index requires the same ordered column list
used to create it. PyPaimon skips composite index files and uses available
single-column indexes or ordinary table scans; composite index building is not
supported in PyPaimon.

## BTree Options

| Option | Default | Description |
Expand All @@ -155,7 +229,7 @@ print(pa_table)
| `btree-index.bloom-filter.enabled` | `false` | Whether to write a Bloom filter to accelerate BTree equality and `IN` lookups. |
| `btree-index.cache-size` | `128 mb` | Cache size used by BTree index readers. |
| `btree-index.high-priority-pool-ratio` | `0.1` | Fraction of `btree-index.cache-size` reserved for high-priority data such as index blocks; the rest caches data blocks. Must be in `[0, 1)`. |
| `btree-index.fallback-scan-max-size` | `256 mb` | Maximum total size of candidate BTree global index files to allow fallback index scans. Set to `0 b` to disable fallback scans. |
| `btree-index.fallback-scan-max-size` | `256 mb` | Maximum total size of candidate BTree global index files to allow fallback index scans, including composite prefix/range scans. Set to `0 b` to disable these scans; complete point probes remain supported. |
| `btree-index.compression` | `none` | Compression algorithm used by BTree index blocks. |
| `btree-index.compression-level` | `1` | Compression level used by codecs that support levels, such as `zstd`. |

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
/*
* 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.paimon.globalindex;

import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.memory.MemorySlice;
import org.apache.paimon.memory.MemorySliceInput;
import org.apache.paimon.memory.MemorySliceOutput;
import org.apache.paimon.types.RowType;

import java.util.Comparator;

/** Length-delimited tuple keys, compared by their typed components with nulls first. */
public class CompositeKeySerializer implements KeySerializer {

private final RowType rowType;

private final KeySerializer[] serializers;
private final InternalRow.FieldGetter[] getters;
private final Comparator<Object>[] comparators;

@SuppressWarnings("unchecked")
public CompositeKeySerializer(RowType type) {
this.rowType = type;
int count = type.getFieldCount();
serializers = new KeySerializer[count];
getters = new InternalRow.FieldGetter[count];
comparators = new Comparator[count];
for (int i = 0; i < count; i++) {
serializers[i] = KeySerializer.create(type.getTypeAt(i));
getters[i] = InternalRow.createFieldGetter(type.getTypeAt(i), i);
comparators[i] = serializers[i].createComparator();
}
}

public RowType rowType() {
return rowType;
}

@Override
public byte[] serialize(Object key) {
InternalRow row = (InternalRow) key;
MemorySliceOutput output = new MemorySliceOutput(32);
for (int i = 0; i < serializers.length; i++) {
Object value = getters[i].getFieldOrNull(row);
if (value == null) {
output.writeInt(-1);
} else {
byte[] bytes = serializers[i].serialize(value);
output.writeInt(bytes.length);
output.writeBytes(bytes);
}
}
return output.toSlice().copyBytes();
}

@Override
public Object deserialize(MemorySlice data) {
MemorySliceInput input = data.toInput();
GenericRow row = new GenericRow(serializers.length);
for (int i = 0; i < serializers.length; i++) {
int length = input.readInt();
if (length >= 0) {
row.setField(i, serializers[i].deserialize(input.readSlice(length)));
} else if (length != -1) {
throw new IllegalArgumentException(
"Invalid composite key component length: " + length);
}
}
if (input.available() != 0) {
throw new IllegalArgumentException("Trailing bytes in composite key");
}
return row;
}

@Override
public Comparator<Object> createComparator() {
return (left, right) -> {
for (int i = 0; i < serializers.length; i++) {
Object a = getters[i].getFieldOrNull((InternalRow) left);
Object b = getters[i].getFieldOrNull((InternalRow) right);
int comparison =
a == null
? (b == null ? 0 : -1)
: b == null ? 1 : comparators[i].compare(a, b);
if (comparison != 0) {
return comparison;
}
}
return 0;
};
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@
import org.apache.paimon.predicate.TopN;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.IOUtils;
import org.apache.paimon.utils.Range;

import javax.annotation.Nullable;

Expand Down Expand Up @@ -306,9 +307,21 @@ public static final class Evaluation {

private final GlobalIndexResult result;
private final Set<Integer> contributingFieldIds;
@Nullable private final List<Range> coveredRanges;

Evaluation(GlobalIndexResult result, Collection<Integer> contributingFieldIds) {
this(result, contributingFieldIds, null);
}

Evaluation(
GlobalIndexResult result,
Collection<Integer> contributingFieldIds,
@Nullable List<Range> coveredRanges) {
this.result = result;
this.coveredRanges =
coveredRanges == null
? null
: Collections.unmodifiableList(new ArrayList<>(coveredRanges));
this.contributingFieldIds =
Collections.unmodifiableSet(new HashSet<>(contributingFieldIds));
}
Expand All @@ -317,6 +330,11 @@ public GlobalIndexResult result() {
return result;
}

@Nullable
public List<Range> coveredRanges() {
return coveredRanges;
}

public Set<Integer> contributingFieldIds() {
return contributingFieldIds;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
import java.io.IOException;
import java.io.Serializable;

/** Extracts zero or more normalized index keys from one source-column value. */
/** Extracts zero or more normalized index keys from one source value or projected tuple. */
public interface GlobalIndexKeyExtractor extends Serializable {

/** Type of the normalized keys emitted by {@link #extract(Object, KeyConsumer)}. */
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import org.apache.paimon.predicate.FullTextSearch;
import org.apache.paimon.predicate.FunctionVisitor;
import org.apache.paimon.predicate.LeafPredicate;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.TopN;
import org.apache.paimon.predicate.VectorSearch;

Expand All @@ -36,6 +37,17 @@
public interface GlobalIndexReader
extends FunctionVisitor<CompletableFuture<Optional<GlobalIndexResult>>>, Closeable {

/** Query an ordered composite key, including prefix ranges and key-only filters. */
default CompletableFuture<Optional<GlobalIndexResult>> visitComposite(Predicate predicate) {
return CompletableFuture.completedFuture(Optional.empty());
}

/** Point lookup of a full composite key, with literals in index column order. */
default CompletableFuture<Optional<GlobalIndexResult>> visitCompositeEqual(
List<Object> literals) {
return CompletableFuture.completedFuture(Optional.empty());
}

@Override
default CompletableFuture<Optional<GlobalIndexResult>> visitIsNaN(FieldRef fieldRef) {
return CompletableFuture.completedFuture(Optional.empty());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -327,7 +327,7 @@ private CompletableFuture<Optional<GlobalIndexResult>> visitParallel(
return visitSelectedFiles(selector.get(), visitor);
}

private CompletableFuture<Optional<GlobalIndexResult>> visitSelectedFiles(
protected CompletableFuture<Optional<GlobalIndexResult>> visitSelectedFiles(
Optional<List<GlobalIndexIOMeta>> selectedOpt,
Function<R, Optional<GlobalIndexResult>> visitor) {
if (!selectedOpt.isPresent()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@
/** A global indexer whose normalized keys must be sorted before they are written. */
public interface SortedGlobalIndexer extends GlobalIndexer {

/** Defines how source-column values are normalized into the keys consumed by the writer. */
/**
* Defines how source values or projected tuples are normalized into the keys consumed by the
* writer.
*/
GlobalIndexKeyExtractor keyExtractor();
}
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@

import org.apache.paimon.compression.BlockCompressionFactory;
import org.apache.paimon.compression.CompressOptions;
import org.apache.paimon.globalindex.CompositeKeySerializer;
import org.apache.paimon.globalindex.GlobalIndexIOMeta;
import org.apache.paimon.globalindex.GlobalIndexKeyExtractor;
import org.apache.paimon.globalindex.GlobalIndexReader;
Expand All @@ -31,13 +32,17 @@
import org.apache.paimon.io.cache.CacheManager;
import org.apache.paimon.options.Options;
import org.apache.paimon.types.DataField;
import org.apache.paimon.types.DataType;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.BloomFilter;
import org.apache.paimon.utils.LazyField;
import org.apache.paimon.utils.Range;

import javax.annotation.Nullable;

import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.ExecutorService;

Expand Down Expand Up @@ -75,8 +80,25 @@ public class BTreeGlobalIndexer implements SortedGlobalIndexer {
private final LazyField<CacheManager> cacheManager;

public BTreeGlobalIndexer(DataField dataField, Options options) {
this.keySerializer = KeySerializer.create(dataField.type());
this.keyExtractor = GlobalIndexKeyExtractor.identity(dataField.type());
this(dataField, Collections.emptyList(), options);
}

public BTreeGlobalIndexer(DataField dataField, List<DataField> extraFields, Options options) {
List<DataField> fields = new ArrayList<>();
fields.add(dataField);
fields.addAll(extraFields);
for (DataField field : fields) {
if (field.type() instanceof RowType) {
throw new UnsupportedOperationException(
"BTree index columns must have scalar types: " + field.name());
}
}
DataType keyType = extraFields.isEmpty() ? dataField.type() : new RowType(fields);
this.keySerializer =
extraFields.isEmpty()
? KeySerializer.create(keyType)
: new CompositeKeySerializer((RowType) keyType);
this.keyExtractor = GlobalIndexKeyExtractor.identity(keyType);
this.options = options;
this.fallbackScanMaxSize =
options.get(BTreeIndexOptions.BTREE_INDEX_FALLBACK_SCAN_MAX_SIZE).getBytes();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,8 @@
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.types.DataField;

import java.util.ArrayList;
import java.util.Collections;
import java.util.List;

/** The {@link GlobalIndexerFactory} for btree index. */
Expand All @@ -45,10 +47,25 @@ public List<GlobalIndexIOMeta> selectFiles(
List<DataField> extraFields,
Predicate predicate,
List<GlobalIndexIOMeta> files) {
if (extraFields != null && !extraFields.isEmpty()) {
List<DataField> fields = new ArrayList<>();
fields.add(indexField);
fields.addAll(extraFields);
return CompositeBTreePredicate.plan(fields, predicate)
.map(plan -> plan.selectFiles(files))
.orElse(files);
}
return SortedFileMetaSelector.selectFiles(
predicate, files, KeySerializer.create(indexField.type()));
}

@Override
public GlobalIndexer create(
DataField indexField, List<DataField> extraFields, Options options) {
return new BTreeGlobalIndexer(
indexField, extraFields == null ? Collections.emptyList() : extraFields, options);
}

@Override
public GlobalIndexer create(DataField dataField, Options options) {
return new BTreeGlobalIndexer(dataField, options);
Expand Down
Loading
Loading