From f0c98102cc012d7eeb6d998e2941ed5e7a8eec41 Mon Sep 17 00:00:00 2001 From: Hongshun Wang Date: Mon, 10 Aug 2026 17:54:24 +0800 Subject: [PATCH 1/2] [flink] Spill cdc log to rocksDB for pk batch read in case of OOM. AI-Contributed/Feature: 0/6 AI-Contributed/UT: 0/0 --- .../batch/KvSnapshotAndLogBatchScanner.java | 110 ++--- .../batch/LakeSnapshotAndLogSplitScanner.java | 130 +++--- .../table/scanner/batch/SortedLogRows.java | 391 ++++++++++++++++++ .../KvSnapshotAndLogBatchScannerTest.java | 181 +++++++- .../scanner/batch/SortedLogRowsTest.java | 348 ++++++++++++++++ .../flink/lake/LakeSplitReaderGenerator.java | 8 +- .../source/reader/FlinkSourceSplitReader.java | 9 +- .../lake/FlussLakeUpsertPartitionReader.scala | 6 +- 8 files changed, 1053 insertions(+), 130 deletions(-) create mode 100644 fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SortedLogRows.java create mode 100644 fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/SortedLogRowsTest.java diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScanner.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScanner.java index 2c6cc1e5c36..899440813fc 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScanner.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScanner.java @@ -23,15 +23,12 @@ import org.apache.fluss.client.table.scanner.ScanRecord; import org.apache.fluss.client.table.scanner.SortMergeReader; import org.apache.fluss.client.table.scanner.log.LogScanner; -import org.apache.fluss.client.table.scanner.log.ScanRecords; import org.apache.fluss.memory.MemorySegment; import org.apache.fluss.metadata.Schema; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.metadata.TableInfo; -import org.apache.fluss.record.ChangeType; import org.apache.fluss.record.LogRecord; import org.apache.fluss.row.InternalRow; -import org.apache.fluss.row.KeyValueRow; import org.apache.fluss.row.encode.KeyEncoder; import org.apache.fluss.types.RowType; import org.apache.fluss.utils.CloseableIterator; @@ -43,10 +40,9 @@ import java.io.UncheckedIOException; import java.time.Duration; import java.util.Comparator; -import java.util.Map; import java.util.NoSuchElementException; -import java.util.TreeMap; +import static org.apache.fluss.client.table.scanner.batch.SortedLogRows.DEFAULT_SPILL_THRESHOLD; import static org.apache.fluss.utils.Preconditions.checkArgument; /** @@ -55,17 +51,14 @@ @Internal public class KvSnapshotAndLogBatchScanner implements BatchScanner { - private final TableBucket tableBucket; - private final long logStoppingOffset; private final int[] keyIndexesInScanRow; @Nullable private final int[] adjustProjectedFields; private final Comparator primaryKeyComparator; - private final Map logRows; + @Nullable private final SortedLogRows logRows; @Nullable private final BatchScanner snapshotScanner; - @Nullable private final LogScanner logScanner; - private boolean logScanFinished; + private boolean logRowsLoaded; private boolean finished; private boolean closed; @@ -79,18 +72,37 @@ public KvSnapshotAndLogBatchScanner( long snapshotId, long logStartingOffset, long logStoppingOffset, - @Nullable int[] projectedFields) { + @Nullable int[] projectedFields, + String scannerTmpDir) { + this( + table, + tableBucket, + snapshotId, + logStartingOffset, + logStoppingOffset, + projectedFields, + scannerTmpDir, + DEFAULT_SPILL_THRESHOLD); + } + + @VisibleForTesting + KvSnapshotAndLogBatchScanner( + Table table, + TableBucket tableBucket, + long snapshotId, + long logStartingOffset, + long logStoppingOffset, + @Nullable int[] projectedFields, + String scannerTmpDir, + int spillThreshold) { checkArgument( table.getTableInfo().hasPrimaryKey(), "KvSnapshotAndLogBatchScanner only supports primary-key tables."); - this.tableBucket = tableBucket; - this.logStoppingOffset = logStoppingOffset; ProjectionPlan projectionPlan = createProjectionPlan(table.getTableInfo(), projectedFields); this.keyIndexesInScanRow = projectionPlan.keyIndexesInScanRow; this.adjustProjectedFields = projectionPlan.adjustProjectedFields; this.primaryKeyComparator = createPrimaryKeyComparator(table.getTableInfo()); - this.logRows = new TreeMap<>(primaryKeyComparator); this.snapshotScanner = snapshotId >= 0 @@ -100,18 +112,30 @@ public KvSnapshotAndLogBatchScanner( : null; boolean emptyLogRange = logStartingOffset >= logStoppingOffset || logStoppingOffset <= 0; - this.logScanFinished = emptyLogRange; + this.logRowsLoaded = emptyLogRange; if (emptyLogRange) { - this.logScanner = null; + this.logRows = null; } else { - this.logScanner = + LogScanner logScanner = table.newScan().project(projectionPlan.scanProjectedFields).createLogScanner(); Long partitionId = tableBucket.getPartitionId(); if (partitionId == null) { - this.logScanner.subscribe(tableBucket.getBucket(), logStartingOffset); + logScanner.subscribe(tableBucket.getBucket(), logStartingOffset); } else { - this.logScanner.subscribe(partitionId, tableBucket.getBucket(), logStartingOffset); + logScanner.subscribe(partitionId, tableBucket.getBucket(), logStartingOffset); } + this.logRows = + new SortedLogRows( + table.getTableInfo() + .getRowType() + .project(projectionPlan.scanProjectedFields), + keyIndexesInScanRow, + primaryKeyComparator, + logScanner, + tableBucket, + logStoppingOffset, + scannerTmpDir, + spillThreshold); } } @@ -122,9 +146,9 @@ public CloseableIterator pollBatch(Duration timeout) throws IOExcep return null; } - // 1. collect the bounded log range into memory before merging. - if (!logScanFinished) { - pollLogRecords(timeout); + // 1. collect the bounded log range before merging. + if (logRows != null && !logRowsLoaded) { + logRowsLoaded = logRows.load(timeout); return CloseableIterator.emptyIterator(); } @@ -138,7 +162,9 @@ public CloseableIterator pollBatch(Duration timeout) throws IOExcep keyIndexesInScanRow, snapshotRecords, primaryKeyComparator, - CloseableIterator.wrap(logRows.values().iterator())); + logRows == null + ? CloseableIterator.emptyIterator() + : logRows.newIterator()); } try { @@ -154,37 +180,6 @@ public CloseableIterator pollBatch(Duration timeout) throws IOExcep return sortMergeIterator; } - private void pollLogRecords(Duration timeout) { - ScanRecords scanRecords = logScanner.poll(timeout); - for (ScanRecord scanRecord : scanRecords.records(tableBucket)) { - long logOffset = scanRecord.logOffset(); - if (logOffset >= logStoppingOffset) { - logScanFinished = true; - break; - } - - reduceLogRecord(scanRecord); - if (logOffset >= logStoppingOffset - 1) { - logScanFinished = true; - break; - } - } - - Long consumedUpToOffset = scanRecords.consumedUpToOffset(tableBucket); - if (consumedUpToOffset != null && consumedUpToOffset >= logStoppingOffset) { - logScanFinished = true; - } - } - - private void reduceLogRecord(ScanRecord scanRecord) { - ChangeType changeType = scanRecord.getChangeType(); - boolean isDelete = - changeType == ChangeType.DELETE || changeType == ChangeType.UPDATE_BEFORE; - KeyValueRow keyValueRow = - new KeyValueRow(keyIndexesInScanRow, scanRecord.getRow(), isDelete); - logRows.put(keyValueRow.keyRow(), keyValueRow); - } - private CloseableIterator createSnapshotRecordIterator(Duration timeout) { if (snapshotScanner == null) { snapshotRecordIterator = CloseableIterator.emptyIterator(); @@ -194,6 +189,11 @@ private CloseableIterator createSnapshotRecordIterator(Duration timeo return snapshotRecordIterator; } + @VisibleForTesting + boolean isLogRowsSpilled() { + return logRows != null && logRows.isSpilled(); + } + @Override public void close() throws IOException { if (closed) { @@ -203,7 +203,7 @@ public void close() throws IOException { IOUtils.closeQuietly(sortMergeIterator); IOUtils.closeQuietly(snapshotRecordIterator); IOUtils.closeQuietly(snapshotScanner); - IOUtils.closeQuietly(logScanner); + IOUtils.closeQuietly(logRows); } private static ProjectionPlan createProjectionPlan( diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java index d8292d11c8c..d67eac2fdd9 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java @@ -19,20 +19,19 @@ package org.apache.fluss.client.table.scanner.batch; import org.apache.fluss.client.table.Table; -import org.apache.fluss.client.table.scanner.ScanRecord; import org.apache.fluss.client.table.scanner.SortMergeReader; import org.apache.fluss.client.table.scanner.log.LogScanner; -import org.apache.fluss.client.table.scanner.log.ScanRecords; import org.apache.fluss.lake.source.LakeSource; import org.apache.fluss.lake.source.LakeSplit; import org.apache.fluss.lake.source.RecordReader; import org.apache.fluss.lake.source.SortedRecordReader; import org.apache.fluss.metadata.TableBucket; -import org.apache.fluss.record.ChangeType; import org.apache.fluss.record.LogRecord; import org.apache.fluss.row.InternalRow; import org.apache.fluss.row.KeyValueRow; +import org.apache.fluss.types.RowType; import org.apache.fluss.utils.CloseableIterator; +import org.apache.fluss.utils.IOUtils; import javax.annotation.Nullable; @@ -40,15 +39,15 @@ import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; -import java.util.Collections; import java.util.Comparator; import java.util.List; -import java.util.Map; -import java.util.TreeMap; + +import static org.apache.fluss.utils.Preconditions.checkState; /** A scanner to merge the lakehouse's snapshot and change log. */ public class LakeSnapshotAndLogSplitScanner implements BatchScanner { + private final TableBucket tableBucket; private final @Nullable List lakeSplits; private Comparator rowComparator; private List> lakeRecordIterators = new ArrayList<>(); @@ -60,15 +59,17 @@ public class LakeSnapshotAndLogSplitScanner implements BatchScanner { // the indexes of primary key in emitted row by lake and fluss private int[] keyIndexesInRow; @Nullable private int[] adjustProjectedFields; + private final RowType scanRowType; + private final String scannerTmpDir; - // the sorted logs in memory, mapping from key -> value - private Map logRows; + private @Nullable SortedLogRows logRows; + private @Nullable LogScanner logScanner; - private final LogScanner logScanner; private final long stoppingOffset; - private boolean logScanFinished; + private boolean logRowsLoaded; private SortMergeReader currentSortMergeReader; + private CloseableIterator currentSortMergeIterator; public LakeSnapshotAndLogSplitScanner( Table table, @@ -77,7 +78,9 @@ public LakeSnapshotAndLogSplitScanner( TableBucket tableBucket, long startingOffset, long stoppingOffset, - @Nullable int[] projectedFields) { + @Nullable int[] projectedFields, + String scannerTmpDir) { + this.tableBucket = tableBucket; this.pkIndexes = table.getTableInfo().getSchema().getPrimaryKeyIndexes(); this.lakeSplits = lakeSplits; this.lakeSource = lakeSource; @@ -90,49 +93,52 @@ public LakeSnapshotAndLogSplitScanner( this.keyIndexesInRow = projectionPlan.keyIndexesInScanRow; this.adjustProjectedFields = projectionPlan.adjustProjectedFields; int[] newProjectedFields = projectionPlan.scanProjectedFields; + this.scanRowType = table.getTableInfo().getRowType().project(newProjectedFields); + this.scannerTmpDir = scannerTmpDir; - this.logScanner = table.newScan().project(newProjectedFields).createLogScanner(); this.lakeSource.withProject( Arrays.stream(newProjectedFields) .mapToObj(field -> new int[] {field}) .toArray(int[][]::new)); - if (tableBucket.getPartitionId() != null) { - this.logScanner.subscribe( - tableBucket.getPartitionId(), tableBucket.getBucket(), startingOffset); + this.logRowsLoaded = startingOffset >= stoppingOffset || stoppingOffset <= 0; + if (logRowsLoaded) { + this.logScanner = null; } else { - this.logScanner.subscribe(tableBucket.getBucket(), startingOffset); + this.logScanner = table.newScan().project(newProjectedFields).createLogScanner(); + if (tableBucket.getPartitionId() != null) { + this.logScanner.subscribe( + tableBucket.getPartitionId(), tableBucket.getBucket(), startingOffset); + } else { + this.logScanner.subscribe(tableBucket.getBucket(), startingOffset); + } } - - this.logScanFinished = startingOffset >= stoppingOffset || stoppingOffset <= 0; } @Nullable @Override public CloseableIterator pollBatch(Duration timeout) throws IOException { - if (logScanFinished) { - initializeLakeRecordIterators(); - if (currentSortMergeReader == null) { - currentSortMergeReader = - new SortMergeReader( - adjustProjectedFields, - keyIndexesInRow, - lakeRecordIterators, - rowComparator, - CloseableIterator.wrap( - logRows == null - ? Collections.emptyIterator() - : logRows.values().iterator())); - } - return currentSortMergeReader.readBatch(); - } else { + if (!logRowsLoaded) { initializeLakeRecordIterators(); - if (logRows == null) { - logRows = new TreeMap<>(rowComparator); - } - pollLogRecords(timeout); - return CloseableIterator.wrap(Collections.emptyIterator()); + initializeLogRows(); + logRowsLoaded = logRows.load(timeout); + return CloseableIterator.emptyIterator(); + } + + initializeLakeRecordIterators(); + if (currentSortMergeReader == null) { + CloseableIterator logRowsIterator = + logRows == null ? CloseableIterator.emptyIterator() : logRows.newIterator(); + currentSortMergeReader = + new SortMergeReader( + adjustProjectedFields, + keyIndexesInRow, + lakeRecordIterators, + rowComparator, + logRowsIterator); } + currentSortMergeIterator = currentSortMergeReader.readBatch(); + return currentSortMergeIterator; } private void initializeLakeRecordIterators() throws IOException { @@ -176,38 +182,32 @@ public boolean requireSortedRecords() { }; } - private void pollLogRecords(Duration timeout) { - ScanRecords scanRecords = logScanner.poll(timeout); - for (ScanRecord scanRecord : scanRecords) { - boolean isDelete = - scanRecord.getChangeType() == ChangeType.DELETE - || scanRecord.getChangeType() == ChangeType.UPDATE_BEFORE; - KeyValueRow keyValueRow = - new KeyValueRow(keyIndexesInRow, scanRecord.getRow(), isDelete); - InternalRow keyRow = keyValueRow.keyRow(); - // upsert the key value row - logRows.put(keyRow, keyValueRow); - if (scanRecord.logOffset() >= stoppingOffset - 1) { - // has reached to the end - logScanFinished = true; - break; - } + private void initializeLogRows() { + if (logRows != null) { + return; } + checkState(logScanner != null, "Log scanner must be initialized."); + logRows = + new SortedLogRows( + scanRowType, + keyIndexesInRow, + rowComparator, + logScanner, + tableBucket, + stoppingOffset, + scannerTmpDir); + logScanner = null; } @Override public void close() throws IOException { - try { - if (logScanner != null) { - logScanner.close(); - } - if (lakeRecordIterators != null) { - for (CloseableIterator iterator : lakeRecordIterators) { - iterator.close(); - } + IOUtils.closeQuietly(currentSortMergeIterator); + IOUtils.closeQuietly(logRows); + IOUtils.closeQuietly(logScanner); + if (lakeRecordIterators != null) { + for (CloseableIterator iterator : lakeRecordIterators) { + IOUtils.closeQuietly(iterator); } - } catch (Exception e) { - throw new IOException("Failed to close resources", e); } } } diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SortedLogRows.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SortedLogRows.java new file mode 100644 index 00000000000..1ab3e111404 --- /dev/null +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SortedLogRows.java @@ -0,0 +1,391 @@ +/* + * 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.fluss.client.table.scanner.batch; + +import org.apache.fluss.annotation.VisibleForTesting; +import org.apache.fluss.client.table.scanner.ScanRecord; +import org.apache.fluss.client.table.scanner.log.LogScanner; +import org.apache.fluss.client.table.scanner.log.ScanRecords; +import org.apache.fluss.metadata.KvFormat; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.record.ChangeType; +import org.apache.fluss.rocksdb.RocksDBHandle; +import org.apache.fluss.rocksdb.RocksIteratorWrapper; +import org.apache.fluss.row.BinaryRow; +import org.apache.fluss.row.InternalRow; +import org.apache.fluss.row.KeyValueRow; +import org.apache.fluss.row.ProjectedRow; +import org.apache.fluss.row.decode.RowDecoder; +import org.apache.fluss.row.serializer.RowSerializer; +import org.apache.fluss.types.DataType; +import org.apache.fluss.types.RowType; +import org.apache.fluss.utils.CloseableIterator; +import org.apache.fluss.utils.FileUtils; +import org.apache.fluss.utils.IOUtils; + +import org.rocksdb.AbstractComparator; +import org.rocksdb.ColumnFamilyOptions; +import org.rocksdb.ComparatorOptions; +import org.rocksdb.DBOptions; +import org.rocksdb.RocksDBException; +import org.rocksdb.RocksIterator; +import org.rocksdb.WriteOptions; + +import javax.annotation.Nullable; +import javax.annotation.concurrent.NotThreadSafe; + +import java.io.Closeable; +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.time.Duration; +import java.util.Arrays; +import java.util.Comparator; +import java.util.Map; +import java.util.NoSuchElementException; +import java.util.TreeMap; + +import static org.apache.fluss.row.BinaryRow.BinaryRowFormat.INDEXED; +import static org.apache.fluss.utils.Preconditions.checkArgument; +import static org.apache.fluss.utils.Preconditions.checkState; + +/** Sorted, deduplicated log rows for snapshot-and-log batch merge. */ +@NotThreadSafe +final class SortedLogRows implements Closeable { + static final int DEFAULT_SPILL_THRESHOLD = 8192; + + private static final byte NORMAL_ROW = 0; + private static final byte TOMBSTONE_ROW = 1; + private static final int VALUE_FLAG_LENGTH = 1; + + private final int[] keyIndexes; + private final Comparator rowComparator; + private final Map memoryRows; + private final ProjectedRow keyProjectedRow; + private final RowSerializer keySerializer; + private final RowSerializer valueSerializer; + private final RowDecoder keyDecoder; + private final RowDecoder valueDecoder; + private final LogScanner logScanner; + private final TableBucket tableBucket; + private final Path scannerTmpDirectory; + private final long stoppingOffset; + private final int spillThreshold; + + private boolean loaded; + private boolean closed; + + @Nullable private Path spillDirectory; + @Nullable private DBOptions dbOptions; + @Nullable private ComparatorOptions comparatorOptions; + @Nullable private InternalRowComparator rocksComparator; + @Nullable private RocksDBHandle rocksDBHandle; + @Nullable private WriteOptions writeOptions; + + SortedLogRows( + RowType rowType, + int[] keyIndexes, + Comparator rowComparator, + LogScanner logScanner, + TableBucket tableBucket, + long stoppingOffset, + String scannerTmpDir) { + this( + rowType, + keyIndexes, + rowComparator, + logScanner, + tableBucket, + stoppingOffset, + scannerTmpDir, + DEFAULT_SPILL_THRESHOLD); + } + + @VisibleForTesting + SortedLogRows( + RowType rowType, + int[] keyIndexes, + Comparator rowComparator, + LogScanner logScanner, + TableBucket tableBucket, + long stoppingOffset, + String scannerTmpDir, + int spillThreshold) { + checkArgument(spillThreshold > 0, "Spill threshold must be positive."); + this.keyIndexes = Arrays.copyOf(keyIndexes, keyIndexes.length); + this.rowComparator = rowComparator; + this.memoryRows = new TreeMap<>(rowComparator); + this.keyProjectedRow = ProjectedRow.from(this.keyIndexes); + + RowType keyRowType = rowType.project(this.keyIndexes); + this.keySerializer = new RowSerializer(toDataTypes(keyRowType), INDEXED); + this.valueSerializer = new RowSerializer(toDataTypes(rowType), INDEXED); + this.keyDecoder = RowDecoder.create(KvFormat.INDEXED, toDataTypes(keyRowType)); + this.valueDecoder = RowDecoder.create(KvFormat.INDEXED, toDataTypes(rowType)); + this.logScanner = logScanner; + this.tableBucket = tableBucket; + this.scannerTmpDirectory = Paths.get(scannerTmpDir); + this.stoppingOffset = stoppingOffset; + this.spillThreshold = spillThreshold; + this.loaded = stoppingOffset <= 0; + } + + boolean load(Duration timeout) throws IOException { + checkNotClosed(); + if (loaded) { + return true; + } + + ScanRecords scanRecords = logScanner.poll(timeout); + for (ScanRecord scanRecord : scanRecords.records(tableBucket)) { + long logOffset = scanRecord.logOffset(); + if (logOffset >= stoppingOffset) { + loaded = true; + break; + } + + put(scanRecord); + if (logOffset >= stoppingOffset - 1) { + loaded = true; + break; + } + } + + Long consumedUpToOffset = scanRecords.consumedUpToOffset(tableBucket); + if (consumedUpToOffset != null && consumedUpToOffset >= stoppingOffset) { + loaded = true; + } + if (loaded) { + IOUtils.closeQuietly(logScanner); + } + return loaded; + } + + private void put(ScanRecord scanRecord) throws IOException { + ChangeType changeType = scanRecord.getChangeType(); + boolean isDelete = + changeType == ChangeType.DELETE || changeType == ChangeType.UPDATE_BEFORE; + put(scanRecord.getRow(), isDelete); + } + + private void put(InternalRow row, boolean isDelete) throws IOException { + checkNotClosed(); + + if (rocksDBHandle == null) { + putToMemory(row, isDelete); + if (memoryRows.size() > spillThreshold) { + spillToRocksDB(); + } + } else { + putToRocksDB(row, isDelete); + } + } + + CloseableIterator newIterator() throws IOException { + checkNotClosed(); + checkState(loaded, "Log rows are not ready for iteration."); + if (rocksDBHandle == null) { + return CloseableIterator.wrap(memoryRows.values().iterator()); + } + + RocksIterator rocksIterator = + rocksDBHandle.getDb().newIterator(rocksDBHandle.getDefaultColumnFamilyHandle()); + RocksIteratorWrapper rocksIteratorWrapper = new RocksIteratorWrapper(rocksIterator); + rocksIteratorWrapper.seekToFirst(); + return new RocksDBLogRowsIterator(rocksIteratorWrapper); + } + + @Override + public void close() { + if (closed) { + return; + } + + closed = true; + IOUtils.closeQuietly(logScanner); + memoryRows.clear(); + IOUtils.closeQuietly(writeOptions); + IOUtils.closeQuietly(rocksDBHandle); + IOUtils.closeQuietly(rocksComparator); + IOUtils.closeQuietly(comparatorOptions); + IOUtils.closeQuietly(dbOptions); + if (spillDirectory != null) { + FileUtils.deleteDirectoryQuietly(spillDirectory.toFile()); + } + } + + @VisibleForTesting + boolean isSpilled() { + return rocksDBHandle != null; + } + + @VisibleForTesting + @Nullable + Path spillDirectory() { + return spillDirectory; + } + + private void putToMemory(InternalRow row, boolean isDelete) { + BinaryRow copiedRow = valueSerializer.toBinaryRow(row).copy(); + KeyValueRow keyValueRow = new KeyValueRow(keyIndexes, copiedRow, isDelete); + memoryRows.put(keyValueRow.keyRow(), keyValueRow); + } + + private void spillToRocksDB() throws IOException { + checkState(rocksDBHandle == null, "Log rows have already been spilled."); + + try { + Files.createDirectories(scannerTmpDirectory); + spillDirectory = Files.createTempDirectory(scannerTmpDirectory, "sorted-log-rows-"); + dbOptions = new DBOptions().setCreateIfMissing(true); + comparatorOptions = new ComparatorOptions().setUseDirectBuffer(false); + rocksComparator = + new InternalRowComparator(comparatorOptions, keyDecoder, rowComparator); + ColumnFamilyOptions columnFamilyOptions = + new ColumnFamilyOptions().setComparator(rocksComparator); + rocksDBHandle = + new RocksDBHandle(spillDirectory.toFile(), dbOptions, columnFamilyOptions); + rocksDBHandle.openDB(); + writeOptions = new WriteOptions().setDisableWAL(true); + + for (KeyValueRow keyValueRow : memoryRows.values()) { + putToRocksDB(keyValueRow.valueRow(), keyValueRow.isDelete()); + } + memoryRows.clear(); + } catch (Exception e) { + close(); + throw new IOException("Failed to spill log rows to RocksDB.", e); + } + } + + private void putToRocksDB(InternalRow row, boolean isDelete) throws IOException { + try { + rocksDBHandle + .getDb() + .put( + rocksDBHandle.getDefaultColumnFamilyHandle(), + writeOptions, + serializeKey(row), + serializeValue(row, isDelete)); + } catch (RocksDBException e) { + throw new IOException("Failed to write log row to RocksDB.", e); + } + } + + private byte[] serializeKey(InternalRow row) { + return toBytes(keySerializer.toBinaryRow(keyProjectedRow.replaceRow(row))); + } + + private byte[] serializeValue(InternalRow row, boolean isDelete) { + byte[] rowBytes = toBytes(valueSerializer.toBinaryRow(row)); + byte[] valueBytes = new byte[VALUE_FLAG_LENGTH + rowBytes.length]; + valueBytes[0] = isDelete ? TOMBSTONE_ROW : NORMAL_ROW; + System.arraycopy(rowBytes, 0, valueBytes, VALUE_FLAG_LENGTH, rowBytes.length); + return valueBytes; + } + + private KeyValueRow deserializeValue(byte[] valueBytes) { + boolean isDelete = valueBytes[0] == TOMBSTONE_ROW; + byte[] rowBytes = Arrays.copyOfRange(valueBytes, VALUE_FLAG_LENGTH, valueBytes.length); + InternalRow valueRow = valueDecoder.decode(rowBytes); + return new KeyValueRow(keyIndexes, valueRow, isDelete); + } + + private void checkNotClosed() { + checkState(!closed, "Sorted log rows has already been closed."); + } + + private static DataType[] toDataTypes(RowType rowType) { + return rowType.getChildren().toArray(new DataType[0]); + } + + private static byte[] toBytes(BinaryRow row) { + byte[] bytes = new byte[row.getSizeInBytes()]; + row.copyTo(bytes, 0); + return bytes; + } + + private static byte[] toBytes(ByteBuffer buffer) { + ByteBuffer duplicate = buffer.duplicate(); + byte[] bytes = new byte[duplicate.remaining()]; + duplicate.get(bytes); + return bytes; + } + + private static class InternalRowComparator extends AbstractComparator { + + private final RowDecoder keyDecoder; + private final Comparator rowComparator; + + InternalRowComparator( + ComparatorOptions comparatorOptions, + RowDecoder keyDecoder, + Comparator rowComparator) { + super(comparatorOptions); + this.keyDecoder = keyDecoder; + this.rowComparator = rowComparator; + } + + @Override + public String name() { + return "fluss-sorted-log-rows-comparator"; + } + + @Override + public int compare(ByteBuffer left, ByteBuffer right) { + return rowComparator.compare( + keyDecoder.decode(toBytes(left)), keyDecoder.decode(toBytes(right))); + } + } + + private class RocksDBLogRowsIterator implements CloseableIterator { + + private final RocksIteratorWrapper rocksIteratorWrapper; + private boolean closed; + + private RocksDBLogRowsIterator(RocksIteratorWrapper rocksIteratorWrapper) { + this.rocksIteratorWrapper = rocksIteratorWrapper; + } + + @Override + public boolean hasNext() { + return !closed && rocksIteratorWrapper.isValid(); + } + + @Override + public KeyValueRow next() { + if (!hasNext()) { + throw new NoSuchElementException(); + } + KeyValueRow keyValueRow = deserializeValue(rocksIteratorWrapper.value()); + rocksIteratorWrapper.next(); + return keyValueRow; + } + + @Override + public void close() { + if (closed) { + return; + } + closed = true; + rocksIteratorWrapper.close(); + } + } +} diff --git a/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScannerTest.java b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScannerTest.java index 9bd6db93278..00ae9f4675a 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScannerTest.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScannerTest.java @@ -17,39 +17,60 @@ package org.apache.fluss.client.table.scanner.batch; +import org.apache.fluss.client.table.Table; +import org.apache.fluss.client.table.scanner.Scan; +import org.apache.fluss.client.table.scanner.ScanRecord; +import org.apache.fluss.client.table.scanner.log.LogScanner; +import org.apache.fluss.client.table.scanner.log.ScanRecords; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.record.ChangeType; import org.apache.fluss.record.LogRecord; import org.apache.fluss.row.GenericRow; import org.apache.fluss.row.InternalRow; import org.apache.fluss.utils.CloseableIterator; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; import javax.annotation.Nullable; +import java.nio.file.Path; import java.time.Duration; +import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.Queue; import java.util.concurrent.LinkedBlockingQueue; +import static org.apache.fluss.record.TestData.DATA1_ROW_TYPE; +import static org.apache.fluss.record.TestData.DATA1_TABLE_ID_PK; +import static org.apache.fluss.record.TestData.DATA1_TABLE_INFO_PK; +import static org.apache.fluss.testutils.DataTestUtils.row; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; /** Test for {@link KvSnapshotAndLogBatchScanner}. */ class KvSnapshotAndLogBatchScannerTest { private static final Duration TIMEOUT = Duration.ofMillis(10); + private @TempDir Path tempDir; + @Test void testSnapshotRecordIteratorPollsAfterEmptyBatch() throws Exception { StubBatchScanner scanner = new StubBatchScanner( Arrays.asList( - Collections.emptyList(), - Collections.singletonList(GenericRow.of(1)), - Collections.emptyList(), - Collections.singletonList(GenericRow.of(2)))); + Collections.emptyList(), + Collections.singletonList(GenericRow.of(1)), + Collections.emptyList(), + Collections.singletonList(GenericRow.of(2)))); KvSnapshotAndLogBatchScanner.SnapshotRecordIterator iterator = new KvSnapshotAndLogBatchScanner.SnapshotRecordIterator(scanner, TIMEOUT); @@ -64,6 +85,113 @@ void testSnapshotRecordIteratorPollsAfterEmptyBatch() throws Exception { assertThat(scanner.pollCount).isEqualTo(5); } + @Test + void testMergeSnapshotAndLogRows() throws Exception { + TableBucket tableBucket = new TableBucket(DATA1_TABLE_ID_PK, 0); + long snapshotId = 1L; + + // snapshot data + StubBatchScanner snapshotScanner = + new StubBatchScanner( + Arrays.asList( + Arrays.asList( + row(DATA1_ROW_TYPE, 0, "snapshot-0"), + row(DATA1_ROW_TYPE, 1, "old-1")), + Arrays.asList( + row(DATA1_ROW_TYPE, 2, "snapshot-2"), + row(DATA1_ROW_TYPE, 3, "delete-me"), + row(DATA1_ROW_TYPE, 5, "snapshot-5")))); + + // cdc log + List logRecords = + Arrays.asList( + logRecord(0, ChangeType.INSERT, 0, "insert-0"), + logRecord(1, ChangeType.UPDATE_BEFORE, 1, "old-1"), + logRecord(2, ChangeType.UPDATE_AFTER, 1, "new-1"), + logRecord(4, ChangeType.DELETE, 2, "delete-2"), + logRecord(4, ChangeType.DELETE, 3, "delete-me"), + logRecord(3, ChangeType.INSERT, 2, "insert-2"), + logRecord(5, ChangeType.INSERT, 4, "insert-4")); + + TestingLogScanner logScanner = + new TestingLogScanner(scanRecords(tableBucket, logRecords, logRecords.size())); + + Table table = mock(Table.class); + Scan scan = mock(Scan.class); + when(table.getTableInfo()).thenReturn(DATA1_TABLE_INFO_PK); + when(table.newScan()).thenReturn(scan); + when(scan.project(any(int[].class))).thenReturn(scan); + when(scan.createBatchScanner(tableBucket, snapshotId)).thenReturn(snapshotScanner); + when(scan.createLogScanner()).thenReturn(logScanner); + + try (KvSnapshotAndLogBatchScanner scanner = + new KvSnapshotAndLogBatchScanner( + table, + tableBucket, + snapshotId, + 0L, + logRecords.size(), + null, + tempDir.toString(), + 3)) { + CloseableIterator loadingRows = scanner.pollBatch(TIMEOUT); + assertThat(loadingRows).isNotNull(); + assertThat(loadingRows.hasNext()).isFalse(); + loadingRows.close(); + + assertThat(scanner.isLogRowsSpilled()).isTrue(); + assertThat(logScanner.closed).isTrue(); + + CloseableIterator mergedRows = scanner.pollBatch(TIMEOUT); + assertThat(mergedRows).isNotNull(); + + List actualRows = collectGenericRows(mergedRows); + assertThat(actualRows) + .containsExactly( + row(DATA1_ROW_TYPE, 0, "insert-0"), + row(DATA1_ROW_TYPE, 1, "new-1"), + row(DATA1_ROW_TYPE, 2, "insert-2"), + row(DATA1_ROW_TYPE, 4, "insert-4"), + row(DATA1_ROW_TYPE, 5, "snapshot-5")); + assertThat(scanner.pollBatch(TIMEOUT)).isNull(); + } + } + + private static LogRecord logRecord(long offset, ChangeType changeType, int key, String value) { + return new ScanRecord(offset, 0L, changeType, row(DATA1_ROW_TYPE, key, value)); + } + + private static ScanRecords scanRecords( + TableBucket tableBucket, List records, long consumedUpToOffset) { + Map> recordsByBucket = new HashMap<>(); + List scanRecords = new ArrayList<>(); + for (LogRecord record : records) { + scanRecords.add((ScanRecord) record); + } + recordsByBucket.put(tableBucket, scanRecords); + + Map consumedUpToOffsets = new HashMap<>(); + consumedUpToOffsets.put(tableBucket, consumedUpToOffset); + return new ScanRecords(recordsByBucket, consumedUpToOffsets); + } + + private static List collectGenericRows(CloseableIterator iterator) { + List rows = new ArrayList<>(); + try { + while (iterator.hasNext()) { + InternalRow internalRow = iterator.next(); + rows.add( + row( + DATA1_ROW_TYPE, + internalRow.getInt(0), + internalRow.getString(1).toString())); + } + } finally { + iterator.close(); + } + return rows; + } + private static class StubBatchScanner implements BatchScanner { private final Queue> batches; @@ -88,4 +216,49 @@ public void close() { // do nothing } } + + private static class TestingLogScanner implements LogScanner { + + private final Queue records = new ArrayDeque<>(); + private boolean closed; + + private TestingLogScanner(ScanRecords... records) { + this.records.addAll(Arrays.asList(records)); + } + + @Override + public ScanRecords poll(Duration timeout) { + return records.isEmpty() ? ScanRecords.EMPTY : records.poll(); + } + + @Override + public void subscribe(int bucket, long offset) { + // do nothing + } + + @Override + public void subscribe(long partitionId, int bucket, long offset) { + // do nothing + } + + @Override + public void unsubscribe(long partitionId, int bucket) { + // do nothing + } + + @Override + public void unsubscribe(int bucket) { + // do nothing + } + + @Override + public void wakeup() { + // do nothing + } + + @Override + public void close() { + closed = true; + } + } } diff --git a/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/SortedLogRowsTest.java b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/SortedLogRowsTest.java new file mode 100644 index 00000000000..d64d617b824 --- /dev/null +++ b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/SortedLogRowsTest.java @@ -0,0 +1,348 @@ +/* + * 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.fluss.client.table.scanner.batch; + +import org.apache.fluss.client.table.scanner.ScanRecord; +import org.apache.fluss.client.table.scanner.log.LogScanner; +import org.apache.fluss.client.table.scanner.log.ScanRecords; +import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.record.ChangeType; +import org.apache.fluss.row.GenericRow; +import org.apache.fluss.row.InternalRow; +import org.apache.fluss.row.KeyValueRow; +import org.apache.fluss.types.DataType; +import org.apache.fluss.types.DataTypes; +import org.apache.fluss.types.RowType; +import org.apache.fluss.utils.CloseableIterator; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.nio.file.Files; +import java.nio.file.Path; +import java.time.Duration; +import java.util.ArrayDeque; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.Comparator; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Queue; +import java.util.stream.Collectors; +import java.util.stream.IntStream; + +import static org.apache.fluss.testutils.DataTestUtils.row; +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link SortedLogRows}. */ +class SortedLogRowsTest { + + private static final RowType ROW_TYPE = + RowType.of( + new DataType[] {DataTypes.INT(), DataTypes.STRING()}, + new String[] {"id", "name"}); + private static final int[] KEY_INDEXES = new int[] {0}; + private static final TableBucket TABLE_BUCKET = new TableBucket(1L, 0); + private static final Duration TIMEOUT = Duration.ofMillis(1); + private static final int TEST_SPILL_THRESHOLD = 4; + private static final Comparator ASCENDING_COMPARATOR = + Comparator.comparingInt(row -> row.getInt(0)); + + private @TempDir Path tempDir; + + @Test + void testDoesNotSpillBeforeThreshold() throws Exception { + List records = new ArrayList<>(); + for (int i = 0; i < TEST_SPILL_THRESHOLD; i++) { + records.add(record(i, row(ROW_TYPE, i, "v" + i))); + } + + try (SortedLogRows logRows = createLogRows(records, TEST_SPILL_THRESHOLD)) { + load(logRows); + + assertThat(logRows.isSpilled()).isFalse(); + List expectedRows = expectedRows(0, TEST_SPILL_THRESHOLD); + assertThat(collectValueRows(logRows)).containsExactlyElementsOf(expectedRows); + } + } + + @Test + void testSpillsWhenDistinctRowsExceedThreshold() throws Exception { + List records = new ArrayList<>(); + for (int i = 0; i <= TEST_SPILL_THRESHOLD; i++) { + records.add(record(i, row(ROW_TYPE, i, "v" + i))); + } + + Path spillDirectory; + try (SortedLogRows logRows = createLogRows(records, TEST_SPILL_THRESHOLD + 1)) { + load(logRows); + + spillDirectory = logRows.spillDirectory(); + assertThat(logRows.isSpilled()).isTrue(); + assertThat(spillDirectory).isNotNull(); + assertThat(spillDirectory).startsWith(tempDir); + assertThat(Files.exists(spillDirectory)).isTrue(); + + List rows = collectRows(logRows); + List expectedRows = expectedRows(0, TEST_SPILL_THRESHOLD + 1); + assertThat(toValueRows(rows)).containsExactlyElementsOf(expectedRows); + assertThat(rows.stream().map(KeyValueRow::isDelete).collect(Collectors.toList())) + .containsOnly(false); + } + + assertThat(Files.exists(spillDirectory)).isFalse(); + } + + @Test + void testDoesNotSpillWhenOnlyRawRecordsExceedThreshold() throws Exception { + List records = new ArrayList<>(); + for (int i = 0; i <= TEST_SPILL_THRESHOLD; i++) { + records.add(record(i, row(ROW_TYPE, 1, "v" + i))); + } + + try (SortedLogRows logRows = createLogRows(records, TEST_SPILL_THRESHOLD + 1)) { + load(logRows); + + assertThat(logRows.isSpilled()).isFalse(); + List expectedRows = + Collections.singletonList(row(ROW_TYPE, 1, "v" + TEST_SPILL_THRESHOLD)); + assertThat(collectValueRows(logRows)).containsExactlyElementsOf(expectedRows); + } + } + + @Test + void testDeduplicatesAndKeepsLastRowInMemory() throws Exception { + List records = + Arrays.asList( + record(0, row(ROW_TYPE, 2, "old")), + record(1, row(ROW_TYPE, 1, "one")), + record(2, row(ROW_TYPE, 2, "new"))); + + try (SortedLogRows logRows = createLogRows(records, 3)) { + load(logRows); + + List expectedRows = + Arrays.asList(row(ROW_TYPE, 1, "one"), row(ROW_TYPE, 2, "new")); + assertThat(collectValueRows(logRows)).containsExactlyElementsOf(expectedRows); + } + } + + @Test + void testDeduplicatesAndKeepsLastRowAfterSpill() throws Exception { + List records = new ArrayList<>(); + for (int i = 0; i <= TEST_SPILL_THRESHOLD; i++) { + records.add(record(i, row(ROW_TYPE, i, "v" + i))); + } + records.add(record(TEST_SPILL_THRESHOLD + 1, row(ROW_TYPE, 1, "new-one"))); + records.add( + record(TEST_SPILL_THRESHOLD + 2, row(ROW_TYPE, TEST_SPILL_THRESHOLD, "new-last"))); + + try (SortedLogRows logRows = createLogRows(records, TEST_SPILL_THRESHOLD + 3)) { + load(logRows); + + assertThat(logRows.isSpilled()).isTrue(); + List expectedRows = + Arrays.asList( + row(ROW_TYPE, 0, "v0"), + row(ROW_TYPE, 1, "new-one"), + row(ROW_TYPE, 2, "v2"), + row(ROW_TYPE, 3, "v3"), + row(ROW_TYPE, 4, "new-last")); + assertThat(collectValueRows(logRows)).containsExactlyElementsOf(expectedRows); + } + } + + @Test + void testSpilledIteratorUsesProvidedComparator() throws Exception { + Comparator descendingComparator = + (row1, row2) -> Integer.compare(row2.getInt(0), row1.getInt(0)); + List records = new ArrayList<>(); + for (int i = 0; i <= TEST_SPILL_THRESHOLD; i++) { + records.add(record(i, row(ROW_TYPE, i, "v" + i))); + } + + try (SortedLogRows logRows = + createLogRows(records, TEST_SPILL_THRESHOLD + 1, descendingComparator)) { + load(logRows); + + assertThat(logRows.isSpilled()).isTrue(); + List expectedReversedRows = expectedRows(0, TEST_SPILL_THRESHOLD + 1); + Collections.reverse(expectedReversedRows); + assertThat(collectValueRows(logRows)).containsExactlyElementsOf(expectedReversedRows); + } + } + + @Test + void testDeleteTombstoneIsPreservedInMemory() throws Exception { + List records = + Arrays.asList( + record(0, row(ROW_TYPE, 1, "old")), + record(1, ChangeType.DELETE, row(ROW_TYPE, 1, "deleted"))); + + try (SortedLogRows logRows = createLogRows(records, 2)) { + load(logRows); + + List rows = collectRows(logRows); + assertThat(rows).hasSize(1); + assertThat(rows.get(0).isDelete()).isTrue(); + assertThat(toGenericRow(rows.get(0).valueRow())).isEqualTo(row(ROW_TYPE, 1, "deleted")); + } + } + + @Test + void testLoadCanBeCalledAcrossPolls() throws Exception { + TestingLogScanner logScanner = + new TestingLogScanner( + scanRecords( + Collections.singletonList(record(0, row(ROW_TYPE, 1, "one"))), 1), + scanRecords( + Collections.singletonList(record(1, row(ROW_TYPE, 2, "two"))), 2)); + + try (SortedLogRows logRows = + new SortedLogRows( + ROW_TYPE, + KEY_INDEXES, + ASCENDING_COMPARATOR, + logScanner, + TABLE_BUCKET, + 2, + tempDir.toString(), + TEST_SPILL_THRESHOLD)) { + assertThat(logRows.load(TIMEOUT)).isFalse(); + assertThat(logRows.load(TIMEOUT)).isTrue(); + List expectedRows = + Arrays.asList(row(ROW_TYPE, 1, "one"), row(ROW_TYPE, 2, "two")); + assertThat(collectValueRows(logRows)).containsExactlyElementsOf(expectedRows); + } + } + + private SortedLogRows createLogRows(List records, long stoppingOffset) { + return createLogRows(records, stoppingOffset, ASCENDING_COMPARATOR); + } + + private SortedLogRows createLogRows( + List records, long stoppingOffset, Comparator rowComparator) { + return new SortedLogRows( + ROW_TYPE, + KEY_INDEXES, + rowComparator, + new TestingLogScanner(scanRecords(records, stoppingOffset)), + TABLE_BUCKET, + stoppingOffset, + tempDir.toString(), + TEST_SPILL_THRESHOLD); + } + + private static ScanRecord record(long offset, InternalRow row) { + return record(offset, ChangeType.INSERT, row); + } + + private static ScanRecord record(long offset, ChangeType changeType, InternalRow row) { + return new ScanRecord(offset, 0L, changeType, row); + } + + private static ScanRecords scanRecords(List records, long consumedUpToOffset) { + Map> recordsByBucket = new HashMap<>(); + recordsByBucket.put(TABLE_BUCKET, records); + Map consumedUpToOffsets = new HashMap<>(); + consumedUpToOffsets.put(TABLE_BUCKET, consumedUpToOffset); + return new ScanRecords(recordsByBucket, consumedUpToOffsets); + } + + private static void load(SortedLogRows logRows) throws Exception { + while (!logRows.load(TIMEOUT)) { + // keep polling until the bounded log range has been fully materialized + } + } + + private static List collectValueRows(SortedLogRows logRows) throws Exception { + return toValueRows(collectRows(logRows)); + } + + private static List toValueRows(List keyValueRows) { + return keyValueRows.stream() + .map(keyValueRow -> toGenericRow(keyValueRow.valueRow())) + .collect(Collectors.toList()); + } + + private static List expectedRows(int startInclusive, int endExclusive) { + return IntStream.range(startInclusive, endExclusive) + .mapToObj(i -> row(ROW_TYPE, i, "v" + i)) + .collect(Collectors.toList()); + } + + private static GenericRow toGenericRow(InternalRow internalRow) { + return row(ROW_TYPE, internalRow.getInt(0), internalRow.getString(1).toString()); + } + + private static List collectRows(SortedLogRows logRows) throws Exception { + List rows = new ArrayList<>(); + try (CloseableIterator iterator = logRows.newIterator()) { + while (iterator.hasNext()) { + rows.add(iterator.next()); + } + } + return rows; + } + + private static class TestingLogScanner implements LogScanner { + + private final Queue records = new ArrayDeque<>(); + + private TestingLogScanner(ScanRecords... records) { + this.records.addAll(Arrays.asList(records)); + } + + @Override + public ScanRecords poll(Duration timeout) { + return records.isEmpty() ? ScanRecords.EMPTY : records.poll(); + } + + @Override + public void subscribe(int bucket, long offset) { + // do nothing + } + + @Override + public void subscribe(long partitionId, int bucket, long offset) { + // do nothing + } + + @Override + public void unsubscribe(long partitionId, int bucket) { + // do nothing + } + + @Override + public void unsubscribe(int bucket) { + // do nothing + } + + @Override + public void wakeup() { + // do nothing + } + + @Override + public void close() { + // do nothing + } + } +} diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/lake/LakeSplitReaderGenerator.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/lake/LakeSplitReaderGenerator.java index 229f1428d91..e8c62ef4ce5 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/lake/LakeSplitReaderGenerator.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/lake/LakeSplitReaderGenerator.java @@ -41,14 +41,17 @@ public class LakeSplitReaderGenerator { private final @Nullable int[] projectedFields; private final @Nullable LakeSource lakeSource; + private final String scannerTmpDir; public LakeSplitReaderGenerator( Table table, @Nullable int[] projectedFields, - @Nullable LakeSource lakeSource) { + @Nullable LakeSource lakeSource, + String scannerTmpDir) { this.table = table; this.projectedFields = projectedFields; this.lakeSource = lakeSource; + this.scannerTmpDir = scannerTmpDir; } public void addSplit(SourceSplitBase split, Queue boundedSplits) { @@ -117,7 +120,8 @@ private BatchScanner getBatchScanner(LakeSnapshotAndFlussLogSplit lakeSplit) { lakeSplit.getTableBucket(), lakeSplit.getStartingOffset(), stoppingOffset, - projectedFields); + projectedFields, + scannerTmpDir); } return lakeBatchScanner; } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java index 50e2143962c..553a3f37483 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java @@ -26,6 +26,7 @@ import org.apache.fluss.client.table.scanner.batch.KvSnapshotAndLogBatchScanner; import org.apache.fluss.client.table.scanner.log.LogScanner; import org.apache.fluss.client.table.scanner.log.ScanRecords; +import org.apache.fluss.config.ConfigOptions; import org.apache.fluss.config.Configuration; import org.apache.fluss.exception.PartitionNotExistException; import org.apache.fluss.flink.lake.LakeSplitReaderGenerator; @@ -103,6 +104,7 @@ public class FlinkSourceSplitReader implements SplitReader lakeSource; @@ -130,6 +132,7 @@ public FlinkSourceSplitReader( new FlinkMetricRegistry(flinkSourceReaderMetrics.getSourceReaderMetricGroup()); this.connection = ConnectionFactory.createConnection(flussConf, flinkMetricRegistry); this.table = connection.getTable(tablePath); + this.scannerTmpDir = flussConf.get(ConfigOptions.CLIENT_SCANNER_IO_TMP_DIR); this.tableId = table.getTableInfo().getTableId(); this.sourceOutputType = sourceOutputType; this.boundedSplits = new ArrayDeque<>(); @@ -244,7 +247,8 @@ public void handleSplitsChanges(SplitsChange splitsChanges) { private LakeSplitReaderGenerator getLakeSplitReader() { if (lakeSplitReaderGenerator == null) { lakeSplitReaderGenerator = - new LakeSplitReaderGenerator(table, projectedFields, checkNotNull(lakeSource)); + new LakeSplitReaderGenerator( + table, projectedFields, checkNotNull(lakeSource), scannerTmpDir); } return lakeSplitReaderGenerator; } @@ -422,7 +426,8 @@ private void checkSnapshotSplitOrStartNext() { "Batch hybrid snapshot log split " + "must have a stopping " + "offset.")), - projectedFields); + projectedFields, + scannerTmpDir); currentBoundedSplitReader = new BoundedSplitReader( batchScanner, hybridSnapshotLogSplit.recordsToSkip()); diff --git a/fluss-spark/fluss-spark-common/src/main/scala/org/apache/fluss/spark/read/lake/FlussLakeUpsertPartitionReader.scala b/fluss-spark/fluss-spark-common/src/main/scala/org/apache/fluss/spark/read/lake/FlussLakeUpsertPartitionReader.scala index 449fbc0626c..a9e4aa3f1e8 100644 --- a/fluss-spark/fluss-spark-common/src/main/scala/org/apache/fluss/spark/read/lake/FlussLakeUpsertPartitionReader.scala +++ b/fluss-spark/fluss-spark-common/src/main/scala/org/apache/fluss/spark/read/lake/FlussLakeUpsertPartitionReader.scala @@ -18,7 +18,7 @@ package org.apache.fluss.spark.read.lake import org.apache.fluss.client.table.scanner.batch.LakeSnapshotAndLogSplitScanner -import org.apache.fluss.config.Configuration +import org.apache.fluss.config.{ConfigOptions, Configuration} import org.apache.fluss.lake.source.{LakeSource, LakeSplit} import org.apache.fluss.metadata.TablePath import org.apache.fluss.row.InternalRow @@ -53,7 +53,9 @@ class FlussLakeUpsertPartitionReader( flussPartition.tableBucket, flussPartition.logStartingOffset, flussPartition.logStoppingOffset, - projection) + projection, + flussConfig.get(ConfigOptions.CLIENT_SCANNER_IO_TMP_DIR) + ) private var mergedIterator: Iterator[InternalRow] = Iterator.empty private var scanFinished = false From 80406398d25e0d17091281b6d59d4a961884b7ca Mon Sep 17 00:00:00 2001 From: Hongshun Wang Date: Mon, 17 Aug 2026 17:42:46 +0800 Subject: [PATCH 2/2] modified based on CR Co-Authored-By: Codex AI-Model: gpt-5.6-sol AI-Contributed/Feature: 137/137 AI-Contributed/UT: 0/0 --- .../batch/KvSnapshotAndLogBatchScanner.java | 22 +-- .../batch/LakeSnapshotAndLogSplitScanner.java | 130 ++++++++--------- .../table/scanner/batch/SortedLogRows.java | 136 ++++++------------ .../scanner/batch/SortedLogRowsTest.java | 30 ++-- .../flink/lake/LakeSplitReaderGenerator.java | 8 +- .../source/reader/FlinkSourceSplitReader.java | 3 +- .../lake/FlussLakeUpsertPartitionReader.scala | 6 +- 7 files changed, 146 insertions(+), 189 deletions(-) diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScanner.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScanner.java index 899440813fc..688e6b4b3be 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScanner.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/KvSnapshotAndLogBatchScanner.java @@ -102,7 +102,8 @@ public KvSnapshotAndLogBatchScanner( ProjectionPlan projectionPlan = createProjectionPlan(table.getTableInfo(), projectedFields); this.keyIndexesInScanRow = projectionPlan.keyIndexesInScanRow; this.adjustProjectedFields = projectionPlan.adjustProjectedFields; - this.primaryKeyComparator = createPrimaryKeyComparator(table.getTableInfo()); + KeyEncoder primaryKeyEncoder = createPrimaryKeyEncoder(table.getTableInfo()); + this.primaryKeyComparator = createPrimaryKeyComparator(primaryKeyEncoder); this.snapshotScanner = snapshotId >= 0 @@ -130,7 +131,7 @@ public KvSnapshotAndLogBatchScanner( .getRowType() .project(projectionPlan.scanProjectedFields), keyIndexesInScanRow, - primaryKeyComparator, + primaryKeyEncoder, logScanner, tableBucket, logStoppingOffset, @@ -214,16 +215,19 @@ private static ProjectionPlan createProjectionPlan( projectedFields); } - private static Comparator createPrimaryKeyComparator(TableInfo tableInfo) { + private static KeyEncoder createPrimaryKeyEncoder(TableInfo tableInfo) { int[] physicalPrimaryKeyIndexes = getPhysicalPrimaryKeyIndexes(tableInfo); RowType primaryKeyRowType = Schema.getKeyRowType(tableInfo.getSchema(), physicalPrimaryKeyIndexes); - KeyEncoder primaryKeyEncoder = - KeyEncoder.ofPrimaryKeyEncoder( - primaryKeyRowType, - tableInfo.getPhysicalPrimaryKeys(), - tableInfo.getTableConfig(), - tableInfo.isDefaultBucketKey()); + return KeyEncoder.ofPrimaryKeyEncoder( + primaryKeyRowType, + tableInfo.getPhysicalPrimaryKeys(), + tableInfo.getTableConfig(), + tableInfo.isDefaultBucketKey()); + } + + private static Comparator createPrimaryKeyComparator( + KeyEncoder primaryKeyEncoder) { return (row1, row2) -> { byte[] key1 = primaryKeyEncoder.encodeKey(row1); byte[] key2 = primaryKeyEncoder.encodeKey(row2); diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java index d67eac2fdd9..d8292d11c8c 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/LakeSnapshotAndLogSplitScanner.java @@ -19,19 +19,20 @@ package org.apache.fluss.client.table.scanner.batch; import org.apache.fluss.client.table.Table; +import org.apache.fluss.client.table.scanner.ScanRecord; import org.apache.fluss.client.table.scanner.SortMergeReader; import org.apache.fluss.client.table.scanner.log.LogScanner; +import org.apache.fluss.client.table.scanner.log.ScanRecords; import org.apache.fluss.lake.source.LakeSource; import org.apache.fluss.lake.source.LakeSplit; import org.apache.fluss.lake.source.RecordReader; import org.apache.fluss.lake.source.SortedRecordReader; import org.apache.fluss.metadata.TableBucket; +import org.apache.fluss.record.ChangeType; import org.apache.fluss.record.LogRecord; import org.apache.fluss.row.InternalRow; import org.apache.fluss.row.KeyValueRow; -import org.apache.fluss.types.RowType; import org.apache.fluss.utils.CloseableIterator; -import org.apache.fluss.utils.IOUtils; import javax.annotation.Nullable; @@ -39,15 +40,15 @@ import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.Comparator; import java.util.List; - -import static org.apache.fluss.utils.Preconditions.checkState; +import java.util.Map; +import java.util.TreeMap; /** A scanner to merge the lakehouse's snapshot and change log. */ public class LakeSnapshotAndLogSplitScanner implements BatchScanner { - private final TableBucket tableBucket; private final @Nullable List lakeSplits; private Comparator rowComparator; private List> lakeRecordIterators = new ArrayList<>(); @@ -59,17 +60,15 @@ public class LakeSnapshotAndLogSplitScanner implements BatchScanner { // the indexes of primary key in emitted row by lake and fluss private int[] keyIndexesInRow; @Nullable private int[] adjustProjectedFields; - private final RowType scanRowType; - private final String scannerTmpDir; - private @Nullable SortedLogRows logRows; - private @Nullable LogScanner logScanner; + // the sorted logs in memory, mapping from key -> value + private Map logRows; + private final LogScanner logScanner; private final long stoppingOffset; - private boolean logRowsLoaded; + private boolean logScanFinished; private SortMergeReader currentSortMergeReader; - private CloseableIterator currentSortMergeIterator; public LakeSnapshotAndLogSplitScanner( Table table, @@ -78,9 +77,7 @@ public LakeSnapshotAndLogSplitScanner( TableBucket tableBucket, long startingOffset, long stoppingOffset, - @Nullable int[] projectedFields, - String scannerTmpDir) { - this.tableBucket = tableBucket; + @Nullable int[] projectedFields) { this.pkIndexes = table.getTableInfo().getSchema().getPrimaryKeyIndexes(); this.lakeSplits = lakeSplits; this.lakeSource = lakeSource; @@ -93,52 +90,49 @@ public LakeSnapshotAndLogSplitScanner( this.keyIndexesInRow = projectionPlan.keyIndexesInScanRow; this.adjustProjectedFields = projectionPlan.adjustProjectedFields; int[] newProjectedFields = projectionPlan.scanProjectedFields; - this.scanRowType = table.getTableInfo().getRowType().project(newProjectedFields); - this.scannerTmpDir = scannerTmpDir; + this.logScanner = table.newScan().project(newProjectedFields).createLogScanner(); this.lakeSource.withProject( Arrays.stream(newProjectedFields) .mapToObj(field -> new int[] {field}) .toArray(int[][]::new)); - this.logRowsLoaded = startingOffset >= stoppingOffset || stoppingOffset <= 0; - if (logRowsLoaded) { - this.logScanner = null; + if (tableBucket.getPartitionId() != null) { + this.logScanner.subscribe( + tableBucket.getPartitionId(), tableBucket.getBucket(), startingOffset); } else { - this.logScanner = table.newScan().project(newProjectedFields).createLogScanner(); - if (tableBucket.getPartitionId() != null) { - this.logScanner.subscribe( - tableBucket.getPartitionId(), tableBucket.getBucket(), startingOffset); - } else { - this.logScanner.subscribe(tableBucket.getBucket(), startingOffset); - } + this.logScanner.subscribe(tableBucket.getBucket(), startingOffset); } + + this.logScanFinished = startingOffset >= stoppingOffset || stoppingOffset <= 0; } @Nullable @Override public CloseableIterator pollBatch(Duration timeout) throws IOException { - if (!logRowsLoaded) { + if (logScanFinished) { initializeLakeRecordIterators(); - initializeLogRows(); - logRowsLoaded = logRows.load(timeout); - return CloseableIterator.emptyIterator(); - } - - initializeLakeRecordIterators(); - if (currentSortMergeReader == null) { - CloseableIterator logRowsIterator = - logRows == null ? CloseableIterator.emptyIterator() : logRows.newIterator(); - currentSortMergeReader = - new SortMergeReader( - adjustProjectedFields, - keyIndexesInRow, - lakeRecordIterators, - rowComparator, - logRowsIterator); + if (currentSortMergeReader == null) { + currentSortMergeReader = + new SortMergeReader( + adjustProjectedFields, + keyIndexesInRow, + lakeRecordIterators, + rowComparator, + CloseableIterator.wrap( + logRows == null + ? Collections.emptyIterator() + : logRows.values().iterator())); + } + return currentSortMergeReader.readBatch(); + } else { + initializeLakeRecordIterators(); + if (logRows == null) { + logRows = new TreeMap<>(rowComparator); + } + pollLogRecords(timeout); + return CloseableIterator.wrap(Collections.emptyIterator()); } - currentSortMergeIterator = currentSortMergeReader.readBatch(); - return currentSortMergeIterator; } private void initializeLakeRecordIterators() throws IOException { @@ -182,32 +176,38 @@ public boolean requireSortedRecords() { }; } - private void initializeLogRows() { - if (logRows != null) { - return; + private void pollLogRecords(Duration timeout) { + ScanRecords scanRecords = logScanner.poll(timeout); + for (ScanRecord scanRecord : scanRecords) { + boolean isDelete = + scanRecord.getChangeType() == ChangeType.DELETE + || scanRecord.getChangeType() == ChangeType.UPDATE_BEFORE; + KeyValueRow keyValueRow = + new KeyValueRow(keyIndexesInRow, scanRecord.getRow(), isDelete); + InternalRow keyRow = keyValueRow.keyRow(); + // upsert the key value row + logRows.put(keyRow, keyValueRow); + if (scanRecord.logOffset() >= stoppingOffset - 1) { + // has reached to the end + logScanFinished = true; + break; + } } - checkState(logScanner != null, "Log scanner must be initialized."); - logRows = - new SortedLogRows( - scanRowType, - keyIndexesInRow, - rowComparator, - logScanner, - tableBucket, - stoppingOffset, - scannerTmpDir); - logScanner = null; } @Override public void close() throws IOException { - IOUtils.closeQuietly(currentSortMergeIterator); - IOUtils.closeQuietly(logRows); - IOUtils.closeQuietly(logScanner); - if (lakeRecordIterators != null) { - for (CloseableIterator iterator : lakeRecordIterators) { - IOUtils.closeQuietly(iterator); + try { + if (logScanner != null) { + logScanner.close(); + } + if (lakeRecordIterators != null) { + for (CloseableIterator iterator : lakeRecordIterators) { + iterator.close(); + } } + } catch (Exception e) { + throw new IOException("Failed to close resources", e); } } } diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SortedLogRows.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SortedLogRows.java index 1ab3e111404..5f5888e72d3 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SortedLogRows.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SortedLogRows.java @@ -21,6 +21,7 @@ import org.apache.fluss.client.table.scanner.ScanRecord; import org.apache.fluss.client.table.scanner.log.LogScanner; import org.apache.fluss.client.table.scanner.log.ScanRecords; +import org.apache.fluss.memory.MemorySegment; import org.apache.fluss.metadata.KvFormat; import org.apache.fluss.metadata.TableBucket; import org.apache.fluss.record.ChangeType; @@ -31,6 +32,7 @@ import org.apache.fluss.row.KeyValueRow; import org.apache.fluss.row.ProjectedRow; import org.apache.fluss.row.decode.RowDecoder; +import org.apache.fluss.row.encode.KeyEncoder; import org.apache.fluss.row.serializer.RowSerializer; import org.apache.fluss.types.DataType; import org.apache.fluss.types.RowType; @@ -38,9 +40,7 @@ import org.apache.fluss.utils.FileUtils; import org.apache.fluss.utils.IOUtils; -import org.rocksdb.AbstractComparator; import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.ComparatorOptions; import org.rocksdb.DBOptions; import org.rocksdb.RocksDBException; import org.rocksdb.RocksIterator; @@ -51,13 +51,11 @@ import java.io.Closeable; import java.io.IOException; -import java.nio.ByteBuffer; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; import java.time.Duration; import java.util.Arrays; -import java.util.Comparator; import java.util.Map; import java.util.NoSuchElementException; import java.util.TreeMap; @@ -66,7 +64,12 @@ import static org.apache.fluss.utils.Preconditions.checkArgument; import static org.apache.fluss.utils.Preconditions.checkState; -/** Sorted, deduplicated log rows for snapshot-and-log batch merge. */ +/** + * Sorted, deduplicated log rows for KV snapshot-and-log batch merge. + * + *

Primary keys are ordered by unsigned lexicographical comparison of their encoded bytes. This + * matches both the KV snapshot scan order and RocksDB's default bytewise comparator. + */ @NotThreadSafe final class SortedLogRows implements Closeable { static final int DEFAULT_SPILL_THRESHOLD = 8192; @@ -76,12 +79,10 @@ final class SortedLogRows implements Closeable { private static final int VALUE_FLAG_LENGTH = 1; private final int[] keyIndexes; - private final Comparator rowComparator; - private final Map memoryRows; + private final KeyEncoder primaryKeyEncoder; + private final Map memoryRows; private final ProjectedRow keyProjectedRow; - private final RowSerializer keySerializer; private final RowSerializer valueSerializer; - private final RowDecoder keyDecoder; private final RowDecoder valueDecoder; private final LogScanner logScanner; private final TableBucket tableBucket; @@ -94,35 +95,14 @@ final class SortedLogRows implements Closeable { @Nullable private Path spillDirectory; @Nullable private DBOptions dbOptions; - @Nullable private ComparatorOptions comparatorOptions; - @Nullable private InternalRowComparator rocksComparator; @Nullable private RocksDBHandle rocksDBHandle; @Nullable private WriteOptions writeOptions; - SortedLogRows( - RowType rowType, - int[] keyIndexes, - Comparator rowComparator, - LogScanner logScanner, - TableBucket tableBucket, - long stoppingOffset, - String scannerTmpDir) { - this( - rowType, - keyIndexes, - rowComparator, - logScanner, - tableBucket, - stoppingOffset, - scannerTmpDir, - DEFAULT_SPILL_THRESHOLD); - } - @VisibleForTesting SortedLogRows( RowType rowType, int[] keyIndexes, - Comparator rowComparator, + KeyEncoder primaryKeyEncoder, LogScanner logScanner, TableBucket tableBucket, long stoppingOffset, @@ -130,14 +110,11 @@ final class SortedLogRows implements Closeable { int spillThreshold) { checkArgument(spillThreshold > 0, "Spill threshold must be positive."); this.keyIndexes = Arrays.copyOf(keyIndexes, keyIndexes.length); - this.rowComparator = rowComparator; - this.memoryRows = new TreeMap<>(rowComparator); + this.primaryKeyEncoder = primaryKeyEncoder; + this.memoryRows = new TreeMap<>(SortedLogRows::compareKeys); this.keyProjectedRow = ProjectedRow.from(this.keyIndexes); - RowType keyRowType = rowType.project(this.keyIndexes); - this.keySerializer = new RowSerializer(toDataTypes(keyRowType), INDEXED); this.valueSerializer = new RowSerializer(toDataTypes(rowType), INDEXED); - this.keyDecoder = RowDecoder.create(KvFormat.INDEXED, toDataTypes(keyRowType)); this.valueDecoder = RowDecoder.create(KvFormat.INDEXED, toDataTypes(rowType)); this.logScanner = logScanner; this.tableBucket = tableBucket; @@ -223,8 +200,6 @@ public void close() { memoryRows.clear(); IOUtils.closeQuietly(writeOptions); IOUtils.closeQuietly(rocksDBHandle); - IOUtils.closeQuietly(rocksComparator); - IOUtils.closeQuietly(comparatorOptions); IOUtils.closeQuietly(dbOptions); if (spillDirectory != null) { FileUtils.deleteDirectoryQuietly(spillDirectory.toFile()); @@ -243,9 +218,16 @@ Path spillDirectory() { } private void putToMemory(InternalRow row, boolean isDelete) { - BinaryRow copiedRow = valueSerializer.toBinaryRow(row).copy(); + BinaryRow copiedRow = copyRow(row); KeyValueRow keyValueRow = new KeyValueRow(keyIndexes, copiedRow, isDelete); - memoryRows.put(keyValueRow.keyRow(), keyValueRow); + memoryRows.put(encodeKey(keyValueRow.keyRow()), keyValueRow); + } + + private BinaryRow copyRow(InternalRow row) { + if (row instanceof BinaryRow) { + return ((BinaryRow) row).copy(); + } + return valueSerializer.toBinaryRow(row).copy(); } private void spillToRocksDB() throws IOException { @@ -255,18 +237,15 @@ private void spillToRocksDB() throws IOException { Files.createDirectories(scannerTmpDirectory); spillDirectory = Files.createTempDirectory(scannerTmpDirectory, "sorted-log-rows-"); dbOptions = new DBOptions().setCreateIfMissing(true); - comparatorOptions = new ComparatorOptions().setUseDirectBuffer(false); - rocksComparator = - new InternalRowComparator(comparatorOptions, keyDecoder, rowComparator); - ColumnFamilyOptions columnFamilyOptions = - new ColumnFamilyOptions().setComparator(rocksComparator); + ColumnFamilyOptions columnFamilyOptions = new ColumnFamilyOptions(); rocksDBHandle = new RocksDBHandle(spillDirectory.toFile(), dbOptions, columnFamilyOptions); rocksDBHandle.openDB(); writeOptions = new WriteOptions().setDisableWAL(true); - for (KeyValueRow keyValueRow : memoryRows.values()) { - putToRocksDB(keyValueRow.valueRow(), keyValueRow.isDelete()); + for (Map.Entry entry : memoryRows.entrySet()) { + KeyValueRow keyValueRow = entry.getValue(); + putToRocksDB(entry.getKey(), keyValueRow.valueRow(), keyValueRow.isDelete()); } memoryRows.clear(); } catch (Exception e) { @@ -276,35 +255,42 @@ private void spillToRocksDB() throws IOException { } private void putToRocksDB(InternalRow row, boolean isDelete) throws IOException { + putToRocksDB(encodeKey(keyProjectedRow.replaceRow(row)), row, isDelete); + } + + private void putToRocksDB(byte[] key, InternalRow row, boolean isDelete) throws IOException { try { rocksDBHandle .getDb() .put( rocksDBHandle.getDefaultColumnFamilyHandle(), writeOptions, - serializeKey(row), + key, serializeValue(row, isDelete)); } catch (RocksDBException e) { throw new IOException("Failed to write log row to RocksDB.", e); } } - private byte[] serializeKey(InternalRow row) { - return toBytes(keySerializer.toBinaryRow(keyProjectedRow.replaceRow(row))); + private byte[] encodeKey(InternalRow keyRow) { + return primaryKeyEncoder.encodeKey(keyRow); } private byte[] serializeValue(InternalRow row, boolean isDelete) { - byte[] rowBytes = toBytes(valueSerializer.toBinaryRow(row)); - byte[] valueBytes = new byte[VALUE_FLAG_LENGTH + rowBytes.length]; + BinaryRow binaryRow = valueSerializer.toBinaryRow(row); + byte[] valueBytes = new byte[VALUE_FLAG_LENGTH + binaryRow.getSizeInBytes()]; valueBytes[0] = isDelete ? TOMBSTONE_ROW : NORMAL_ROW; - System.arraycopy(rowBytes, 0, valueBytes, VALUE_FLAG_LENGTH, rowBytes.length); + binaryRow.copyTo(valueBytes, VALUE_FLAG_LENGTH); return valueBytes; } private KeyValueRow deserializeValue(byte[] valueBytes) { boolean isDelete = valueBytes[0] == TOMBSTONE_ROW; - byte[] rowBytes = Arrays.copyOfRange(valueBytes, VALUE_FLAG_LENGTH, valueBytes.length); - InternalRow valueRow = valueDecoder.decode(rowBytes); + InternalRow valueRow = + valueDecoder.decode( + MemorySegment.wrap(valueBytes), + VALUE_FLAG_LENGTH, + valueBytes.length - VALUE_FLAG_LENGTH); return new KeyValueRow(keyIndexes, valueRow, isDelete); } @@ -316,43 +302,9 @@ private static DataType[] toDataTypes(RowType rowType) { return rowType.getChildren().toArray(new DataType[0]); } - private static byte[] toBytes(BinaryRow row) { - byte[] bytes = new byte[row.getSizeInBytes()]; - row.copyTo(bytes, 0); - return bytes; - } - - private static byte[] toBytes(ByteBuffer buffer) { - ByteBuffer duplicate = buffer.duplicate(); - byte[] bytes = new byte[duplicate.remaining()]; - duplicate.get(bytes); - return bytes; - } - - private static class InternalRowComparator extends AbstractComparator { - - private final RowDecoder keyDecoder; - private final Comparator rowComparator; - - InternalRowComparator( - ComparatorOptions comparatorOptions, - RowDecoder keyDecoder, - Comparator rowComparator) { - super(comparatorOptions); - this.keyDecoder = keyDecoder; - this.rowComparator = rowComparator; - } - - @Override - public String name() { - return "fluss-sorted-log-rows-comparator"; - } - - @Override - public int compare(ByteBuffer left, ByteBuffer right) { - return rowComparator.compare( - keyDecoder.decode(toBytes(left)), keyDecoder.decode(toBytes(right))); - } + private static int compareKeys(byte[] left, byte[] right) { + return MemorySegment.wrap(left) + .compare(MemorySegment.wrap(right), 0, 0, left.length, right.length); } private class RocksDBLogRowsIterator implements CloseableIterator { diff --git a/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/SortedLogRowsTest.java b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/SortedLogRowsTest.java index d64d617b824..7fea0e95f6c 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/SortedLogRowsTest.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/table/scanner/batch/SortedLogRowsTest.java @@ -25,6 +25,7 @@ import org.apache.fluss.row.GenericRow; import org.apache.fluss.row.InternalRow; import org.apache.fluss.row.KeyValueRow; +import org.apache.fluss.row.encode.KeyEncoder; import org.apache.fluss.types.DataType; import org.apache.fluss.types.DataTypes; import org.apache.fluss.types.RowType; @@ -40,7 +41,6 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; -import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -62,8 +62,7 @@ class SortedLogRowsTest { private static final TableBucket TABLE_BUCKET = new TableBucket(1L, 0); private static final Duration TIMEOUT = Duration.ofMillis(1); private static final int TEST_SPILL_THRESHOLD = 4; - private static final Comparator ASCENDING_COMPARATOR = - Comparator.comparingInt(row -> row.getInt(0)); + private static final KeyEncoder ASCENDING_KEY_ENCODER = row -> encodeSortableInt(row.getInt(0)); private @TempDir Path tempDir; @@ -170,16 +169,15 @@ void testDeduplicatesAndKeepsLastRowAfterSpill() throws Exception { } @Test - void testSpilledIteratorUsesProvidedComparator() throws Exception { - Comparator descendingComparator = - (row1, row2) -> Integer.compare(row2.getInt(0), row1.getInt(0)); + void testSpilledIteratorUsesEncodedKeyOrder() throws Exception { + KeyEncoder descendingKeyEncoder = row -> encodeSortableInt(-row.getInt(0)); List records = new ArrayList<>(); for (int i = 0; i <= TEST_SPILL_THRESHOLD; i++) { records.add(record(i, row(ROW_TYPE, i, "v" + i))); } try (SortedLogRows logRows = - createLogRows(records, TEST_SPILL_THRESHOLD + 1, descendingComparator)) { + createLogRows(records, TEST_SPILL_THRESHOLD + 1, descendingKeyEncoder)) { load(logRows); assertThat(logRows.isSpilled()).isTrue(); @@ -219,7 +217,7 @@ void testLoadCanBeCalledAcrossPolls() throws Exception { new SortedLogRows( ROW_TYPE, KEY_INDEXES, - ASCENDING_COMPARATOR, + ASCENDING_KEY_ENCODER, logScanner, TABLE_BUCKET, 2, @@ -234,15 +232,15 @@ void testLoadCanBeCalledAcrossPolls() throws Exception { } private SortedLogRows createLogRows(List records, long stoppingOffset) { - return createLogRows(records, stoppingOffset, ASCENDING_COMPARATOR); + return createLogRows(records, stoppingOffset, ASCENDING_KEY_ENCODER); } private SortedLogRows createLogRows( - List records, long stoppingOffset, Comparator rowComparator) { + List records, long stoppingOffset, KeyEncoder primaryKeyEncoder) { return new SortedLogRows( ROW_TYPE, KEY_INDEXES, - rowComparator, + primaryKeyEncoder, new TestingLogScanner(scanRecords(records, stoppingOffset)), TABLE_BUCKET, stoppingOffset, @@ -250,6 +248,16 @@ private SortedLogRows createLogRows( TEST_SPILL_THRESHOLD); } + private static byte[] encodeSortableInt(int value) { + int normalized = value ^ Integer.MIN_VALUE; + return new byte[] { + (byte) (normalized >>> 24), + (byte) (normalized >>> 16), + (byte) (normalized >>> 8), + (byte) normalized + }; + } + private static ScanRecord record(long offset, InternalRow row) { return record(offset, ChangeType.INSERT, row); } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/lake/LakeSplitReaderGenerator.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/lake/LakeSplitReaderGenerator.java index e8c62ef4ce5..229f1428d91 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/lake/LakeSplitReaderGenerator.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/lake/LakeSplitReaderGenerator.java @@ -41,17 +41,14 @@ public class LakeSplitReaderGenerator { private final @Nullable int[] projectedFields; private final @Nullable LakeSource lakeSource; - private final String scannerTmpDir; public LakeSplitReaderGenerator( Table table, @Nullable int[] projectedFields, - @Nullable LakeSource lakeSource, - String scannerTmpDir) { + @Nullable LakeSource lakeSource) { this.table = table; this.projectedFields = projectedFields; this.lakeSource = lakeSource; - this.scannerTmpDir = scannerTmpDir; } public void addSplit(SourceSplitBase split, Queue boundedSplits) { @@ -120,8 +117,7 @@ private BatchScanner getBatchScanner(LakeSnapshotAndFlussLogSplit lakeSplit) { lakeSplit.getTableBucket(), lakeSplit.getStartingOffset(), stoppingOffset, - projectedFields, - scannerTmpDir); + projectedFields); } return lakeBatchScanner; } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java index 553a3f37483..3656cb0209c 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/source/reader/FlinkSourceSplitReader.java @@ -247,8 +247,7 @@ public void handleSplitsChanges(SplitsChange splitsChanges) { private LakeSplitReaderGenerator getLakeSplitReader() { if (lakeSplitReaderGenerator == null) { lakeSplitReaderGenerator = - new LakeSplitReaderGenerator( - table, projectedFields, checkNotNull(lakeSource), scannerTmpDir); + new LakeSplitReaderGenerator(table, projectedFields, checkNotNull(lakeSource)); } return lakeSplitReaderGenerator; } diff --git a/fluss-spark/fluss-spark-common/src/main/scala/org/apache/fluss/spark/read/lake/FlussLakeUpsertPartitionReader.scala b/fluss-spark/fluss-spark-common/src/main/scala/org/apache/fluss/spark/read/lake/FlussLakeUpsertPartitionReader.scala index a9e4aa3f1e8..449fbc0626c 100644 --- a/fluss-spark/fluss-spark-common/src/main/scala/org/apache/fluss/spark/read/lake/FlussLakeUpsertPartitionReader.scala +++ b/fluss-spark/fluss-spark-common/src/main/scala/org/apache/fluss/spark/read/lake/FlussLakeUpsertPartitionReader.scala @@ -18,7 +18,7 @@ package org.apache.fluss.spark.read.lake import org.apache.fluss.client.table.scanner.batch.LakeSnapshotAndLogSplitScanner -import org.apache.fluss.config.{ConfigOptions, Configuration} +import org.apache.fluss.config.Configuration import org.apache.fluss.lake.source.{LakeSource, LakeSplit} import org.apache.fluss.metadata.TablePath import org.apache.fluss.row.InternalRow @@ -53,9 +53,7 @@ class FlussLakeUpsertPartitionReader( flussPartition.tableBucket, flussPartition.logStartingOffset, flussPartition.logStoppingOffset, - projection, - flussConfig.get(ConfigOptions.CLIENT_SCANNER_IO_TMP_DIR) - ) + projection) private var mergedIterator: Iterator[InternalRow] = Iterator.empty private var scanFinished = false