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
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
/*
* 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;
if (row.getFieldCount() != serializers.length) {
throw new IllegalArgumentException(
"Expected "
+ serializers.length
+ " composite key fields, but got "
+ row.getFieldCount());
}
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 @@ -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 @@ -36,6 +36,12 @@
public interface GlobalIndexReader
extends FunctionVisitor<CompletableFuture<Optional<GlobalIndexResult>>>, Closeable {

/** 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 @@ -49,14 +49,7 @@ GlobalIndexReader createReader(
@Nullable List<Range> rowRanges,
ExecutorService executor);

static GlobalIndexer create(String type, DataField indexField, Options options) {
GlobalIndexerFactory globalIndexerFactory = GlobalIndexerFactoryUtils.load(type);
return globalIndexerFactory.create(indexField, options);
}

static GlobalIndexer create(
String type, DataField indexField, List<DataField> extraFields, Options options) {
GlobalIndexerFactory globalIndexerFactory = GlobalIndexerFactoryUtils.load(type);
return globalIndexerFactory.create(indexField, extraFields, options);
static GlobalIndexer create(String type, List<DataField> indexFields, Options options) {
return GlobalIndexerFactoryUtils.load(type).create(indexFields, options);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -39,28 +39,10 @@ default boolean supportsFullTextSearch() {
* match; the default keeps all files when no safe pruning is available.
*/
default List<GlobalIndexIOMeta> selectFiles(
DataField indexField,
List<DataField> extraFields,
Predicate predicate,
List<GlobalIndexIOMeta> files) {
List<DataField> indexFields, Predicate predicate, List<GlobalIndexIOMeta> files) {
return files;
}

GlobalIndexer create(DataField indexField, Options options);

/**
* Creates an indexer over a primary column plus optional extra columns. {@code indexField} is
* the primary column; {@code extraFields} holds the remaining columns and is empty for a
* single-column index.
*/
default GlobalIndexer create(
DataField indexField, List<DataField> extraFields, Options options) {
if (extraFields != null && !extraFields.isEmpty()) {
throw new UnsupportedOperationException(
String.format(
"Index type '%s' does not support multi-column index, got extra columns: %s",
identifier(), extraFields));
}
return create(indexField, options);
}
/** Creates an indexer over a non-empty list of columns in index order. */
GlobalIndexer create(List<DataField> indexFields, Options options);
}
Original file line number Diff line number Diff line change
Expand Up @@ -61,13 +61,10 @@ public static GlobalIndexerFactory load(String type) {
/** Keep unrecognized index types for the reader to resolve at execution time. */
public static List<GlobalIndexIOMeta> selectFiles(
String type,
DataField indexField,
List<DataField> extraFields,
List<DataField> indexFields,
Predicate predicate,
List<GlobalIndexIOMeta> files) {
GlobalIndexerFactory factory = factories.get(type);
return factory == null
? files
: factory.selectFiles(indexField, extraFields, predicate, files);
return factory == null ? files : factory.selectFiles(indexFields, predicate, files);
}
}
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 @@ -41,16 +41,18 @@ public String identifier() {

@Override
public List<GlobalIndexIOMeta> selectFiles(
DataField indexField,
List<DataField> extraFields,
Predicate predicate,
List<GlobalIndexIOMeta> files) {
List<DataField> indexFields, Predicate predicate, List<GlobalIndexIOMeta> files) {
return SortedFileMetaSelector.selectFiles(
predicate, files, KeySerializer.create(indexField.type()));
predicate, files, KeySerializer.create(indexFields.get(0).type()));
}

@Override
public GlobalIndexer create(DataField dataField, Options options) {
public GlobalIndexer create(List<DataField> indexFields, Options options) {
if (indexFields.size() != 1) {
throw new UnsupportedOperationException(
"Index type '" + identifier() + "' requires exactly one index field.");
}
DataField dataField = indexFields.get(0);
return new BitmapGlobalIndexer(dataField, options);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@
import org.apache.paimon.options.Options;
import org.apache.paimon.types.DataField;

import java.util.List;

/** Factory for bitmap-backed multivalue indexes on array columns. */
public class MultiValueGlobalIndexerFactory implements GlobalIndexerFactory {

Expand All @@ -34,7 +36,12 @@ public String identifier() {
}

@Override
public GlobalIndexer create(DataField indexField, Options options) {
public GlobalIndexer create(List<DataField> indexFields, Options options) {
if (indexFields.size() != 1) {
throw new UnsupportedOperationException(
"Index type '" + identifier() + "' requires exactly one index field.");
}
DataField indexField = indexFields.get(0);
return new MultiValueGlobalIndexer(indexField, options);
}
}
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,6 +32,8 @@
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;
Expand All @@ -41,6 +44,8 @@
import java.util.List;
import java.util.concurrent.ExecutorService;

import static org.apache.paimon.utils.Preconditions.checkArgument;

/**
* The {@link GlobalIndexer} for btree index. We do not build a B-tree directly in memory, instead,
* we form a logical B-tree via multi-level metadata over SST files that store the actual data, as
Expand Down Expand Up @@ -74,9 +79,20 @@ public class BTreeGlobalIndexer implements SortedGlobalIndexer {
private final long fallbackScanMaxSize;
private final LazyField<CacheManager> cacheManager;

public BTreeGlobalIndexer(DataField dataField, Options options) {
this.keySerializer = KeySerializer.create(dataField.type());
this.keyExtractor = GlobalIndexKeyExtractor.identity(dataField.type());
public BTreeGlobalIndexer(List<DataField> fields, Options options) {
checkArgument(!fields.isEmpty(), "BTree index requires at least one field.");
for (DataField field : fields) {
if (field.type() instanceof RowType) {
throw new UnsupportedOperationException(
"BTree index columns must have scalar types: " + field.name());
}
}
DataType keyType = fields.size() == 1 ? fields.get(0).type() : new RowType(fields);
this.keySerializer =
fields.size() == 1
? 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 @@ -41,16 +41,17 @@ public String identifier() {

@Override
public List<GlobalIndexIOMeta> selectFiles(
DataField indexField,
List<DataField> extraFields,
Predicate predicate,
List<GlobalIndexIOMeta> files) {
List<DataField> indexFields, Predicate predicate, List<GlobalIndexIOMeta> files) {
if (indexFields.size() > 1) {
// Scalar predicates cannot safely prune tuple metadata.
return files;
}
return SortedFileMetaSelector.selectFiles(
predicate, files, KeySerializer.create(indexField.type()));
predicate, files, KeySerializer.create(indexFields.get(0).type()));
}

@Override
public GlobalIndexer create(DataField dataField, Options options) {
return new BTreeGlobalIndexer(dataField, options);
public GlobalIndexer create(List<DataField> indexFields, Options options) {
return new BTreeGlobalIndexer(indexFields, options);
}
}
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.fs.PositionOutputStream;
import org.apache.paimon.globalindex.CompositeKeySerializer;
import org.apache.paimon.globalindex.GlobalIndexSingleColumnWriter;
import org.apache.paimon.globalindex.KeySerializer;
import org.apache.paimon.globalindex.ResultEntry;
Expand Down Expand Up @@ -148,19 +149,27 @@ public void write(@Nullable Object key, long rowId) {
return;
}

if (lastKey != null && comparator.compare(key, lastKey) != 0) {
boolean differentKey = lastKey == null || comparator.compare(key, lastKey) != 0;
if (lastKey != null && differentKey) {
try {
flush();
} catch (IOException e) {
throw new RuntimeException("Error in writing btree index files.", e);
}
}
lastKey = key;
if (differentKey || !(keySerializer instanceof CompositeKeySerializer)) {
// Sorted engine iterators can reuse the row backing a tuple key.
lastKey =
keySerializer instanceof CompositeKeySerializer
? keySerializer.deserialize(
MemorySlice.wrap(keySerializer.serialize(key)))
: key;
}
currentRowIds.add(rowId);

// update stats
if (firstKey == null) {
firstKey = key;
firstKey = lastKey;
}
}

Expand Down
Loading
Loading