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
Expand Up @@ -20,87 +20,56 @@

import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.serializer.RowCompactedSerializer;
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. */
/** Compacted row keys, compared by their typed components with nulls first. */
public class CompositeKeySerializer implements KeySerializer {

private final RowType rowType;

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

@SuppressWarnings("unchecked")
public CompositeKeySerializer(RowType type) {
this.rowType = type;
serializer = ThreadLocal.withInitial(() -> new RowCompactedSerializer(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();
comparators[i] = KeySerializer.create(type.getTypeAt(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);
// Canonicalize RowKind and NaN payloads so equal keys have the same Bloom hash.
GenericRow row = new GenericRow(getters.length);
for (int i = 0; i < getters.length; i++) {
Object value = getters[i].getFieldOrNull((InternalRow) key);
if (value instanceof Float && Float.isNaN((Float) value)) {
value = Float.NaN;
} else if (value instanceof Double && Double.isNaN((Double) value)) {
value = Double.NaN;
}
row.setField(i, value);
}
return output.toSlice().copyBytes();
return serializer.get().serializeToBytes(row);
}

@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;
return serializer.get().deserialize(data.copyBytes());
}

@Override
public Comparator<Object> createComparator() {
return (left, right) -> {
for (int i = 0; i < serializers.length; i++) {
for (int i = 0; i < getters.length; i++) {
Object a = getters[i].getFieldOrNull((InternalRow) left);
Object b = getters[i].getFieldOrNull((InternalRow) right);
int comparison =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,8 @@ CompletableFuture<Optional<GlobalIndexEvaluator.Evaluation>> evaluate(
Optional.of(
new GlobalIndexEvaluator.Evaluation(
GlobalIndexResult.createEmpty(),
emptyResultFields)));
emptyResultFields,
null)));
}
}
List<ContainsGroup> groupedContains = new ArrayList<>(groups.values());
Expand Down Expand Up @@ -169,7 +170,8 @@ CompletableFuture<Optional<GlobalIndexEvaluator.Evaluation>> evaluate(
.Evaluation(
GlobalIndexResult
.createEmpty(),
emptyResultFields)));
emptyResultFields,
null)));
}
}

Expand Down Expand Up @@ -269,7 +271,8 @@ private CompletableFuture<Optional<GlobalIndexEvaluator.Evaluation>> visitContai
value ->
new GlobalIndexEvaluator.Evaluation(
value,
Collections.singleton(group.fieldId))));
Collections.singleton(group.fieldId),
null)));
}

private CompletableFuture<Optional<GlobalIndexResult>> intersectReaderResults(
Expand Down
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 @@ -188,7 +189,8 @@ private CompletableFuture<Optional<Evaluation>> visitFieldAsync(
}
return compoundResult.map(
result ->
new Evaluation(result, Collections.singleton(fieldId)));
new Evaluation(
result, Collections.singleton(fieldId), null));
});
}

Expand Down Expand Up @@ -278,7 +280,7 @@ private Optional<Evaluation> combineResults(
compoundResult = compoundResult.or(child.get().result());
contributingFieldIds.addAll(child.get().contributingFieldIds());
}
return Optional.of(new Evaluation(compoundResult, contributingFieldIds));
return Optional.of(new Evaluation(compoundResult, contributingFieldIds, null));
} else {
Optional<GlobalIndexResult> compoundResult = Optional.empty();
for (Optional<Evaluation> child : results) {
Expand All @@ -295,7 +297,7 @@ private Optional<Evaluation> combineResults(
break;
}
}
return compoundResult.map(result -> new Evaluation(result, contributingFieldIds));
return compoundResult.map(result -> new Evaluation(result, contributingFieldIds, null));
}
}

Expand All @@ -306,9 +308,17 @@ 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) {
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 +327,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 @@ -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,9 +37,8 @@
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) {
/** Evaluate a composite predicate. An empty Optional means the reader cannot use the index. */
default CompletableFuture<Optional<GlobalIndexResult>> visitComposite(Predicate predicate) {
return CompletableFuture.completedFuture(Optional.empty());
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
import javax.annotation.Nullable;

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

Expand Down Expand Up @@ -74,12 +75,14 @@ public class BTreeGlobalIndexer implements SortedGlobalIndexer {
private static final double BLOOM_FILTER_FPP = 0.05;

private final KeySerializer keySerializer;
private final List<DataField> fields;
private final GlobalIndexKeyExtractor keyExtractor;
private final Options options;
private final long fallbackScanMaxSize;
private final LazyField<CacheManager> cacheManager;

public BTreeGlobalIndexer(List<DataField> fields, Options options) {
this.fields = new ArrayList<>(fields);
checkArgument(!fields.isEmpty(), "BTree index requires at least one field.");
for (DataField field : fields) {
if (field.type() instanceof RowType) {
Expand Down Expand Up @@ -140,6 +143,7 @@ public GlobalIndexReader createReader(
ExecutorService executor) {
return new LazyFilteredBTreeReader(
files,
fields,
keySerializer,
fileReader,
cacheManager.get(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,16 +18,23 @@

package org.apache.paimon.globalindex.btree;

import org.apache.paimon.data.GenericRow;
import org.apache.paimon.globalindex.CompositeKeySerializer;
import org.apache.paimon.globalindex.GlobalIndexIOMeta;
import org.apache.paimon.globalindex.GlobalIndexer;
import org.apache.paimon.globalindex.GlobalIndexerFactory;
import org.apache.paimon.globalindex.KeySerializer;
import org.apache.paimon.globalindex.SortedFileMetaSelector;
import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.FieldRef;
import org.apache.paimon.predicate.LeafPredicate;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.types.DataField;
import org.apache.paimon.types.RowType;

import java.util.Collections;
import java.util.List;
import java.util.Optional;

/** The {@link GlobalIndexerFactory} for btree index. */
public class BTreeGlobalIndexerFactory implements GlobalIndexerFactory {
Expand All @@ -43,8 +50,25 @@ public String identifier() {
public List<GlobalIndexIOMeta> selectFiles(
List<DataField> indexFields, Predicate predicate, List<GlobalIndexIOMeta> files) {
if (indexFields.size() > 1) {
// Scalar predicates cannot safely prune tuple metadata.
return files;
Optional<List<LeafPredicate>> matched =
CompositeBTreePredicate.match(indexFields, predicate);
if (!matched.isPresent()) {
return files;
}
if (CompositeBTreePredicate.isContradictory(indexFields, predicate)) {
return Collections.emptyList();
}
if (files.stream().anyMatch(file -> file.metadata() == null)) {
return files;
}
RowType keyType = new RowType(indexFields);
KeySerializer serializer = new CompositeKeySerializer(keyType);
Object[] values = matched.get().stream().map(leaf -> leaf.literals().get(0)).toArray();
return new SortedFileMetaSelector(files, serializer)
.visitEqual(
new FieldRef(0, indexFields.get(0).name(), keyType),
GenericRow.of(values))
.orElse(files);
}
return SortedFileMetaSelector.selectFiles(
predicate, files, KeySerializer.create(indexFields.get(0).type()));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@

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 @@ -149,27 +148,19 @@ public void write(@Nullable Object key, long rowId) {
return;
}

boolean differentKey = lastKey == null || comparator.compare(key, lastKey) != 0;
if (lastKey != null && differentKey) {
if (lastKey != null && comparator.compare(key, lastKey) != 0) {
try {
flush();
} catch (IOException e) {
throw new RuntimeException("Error in writing btree index files.", e);
}
}
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;
}
lastKey = key;
currentRowIds.add(rowId);

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

Expand Down
Loading
Loading