This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new ee771269f6 [core] Replace RocksDB-backed caches with LocalKvDb (#8881)
ee771269f6 is described below
commit ee771269f69aabb19f8b9e42a4d451a43394a26c
Author: Jingsong Lee <[email protected]>
AuthorDate: Tue Jul 28 15:18:46 2026 +0800
[core] Replace RocksDB-backed caches with LocalKvDb (#8881)
---
docs/generated/flink_connector_configuration.html | 2 +-
.../benchmark/lookup/LocalKvDbBenchmark.java | 330 ++++++++++++++++++++-
.../apache/paimon/lookup/sort/db/LocalKvDb.java | 231 +++++++++++----
.../lookup/sort/db/RecordCombiningWriter.java | 8 +-
.../paimon/lookup/sort/db/LocalKvDbTest.java | 112 +++----
.../paimon/crosspartition/GlobalIndexAssigner.java | 46 +--
.../lookup/local/LocalKvListMergeOperator.java | 5 -
.../paimon/lookup/local/LocalKvListValueCodec.java | 91 +++---
.../crosspartition/GlobalIndexAssignerTest.java | 13 +-
.../lookup/local/LocalKvStateFactoryTest.java | 48 +++
.../apache/paimon/flink/FlinkConnectorOptions.java | 2 +-
.../paimon/flink/lookup/FullCacheLookupTable.java | 14 +-
.../flink/lookup/NoPrimaryKeyLookupTable.java | 4 +-
.../paimon/flink/lookup/PrimaryKeyLookupTable.java | 4 +-
.../flink/sink/index/GlobalDynamicBucketSink.java | 2 +-
.../sink/index/GlobalIndexAssignerOperator.java | 1 +
.../paimon/flink/CrossPartitionTableITCase.java | 1 -
.../paimon/flink/lookup/LookupTableTest.java | 2 +-
.../paimon/spark/commands/BucketProcessor.scala | 1 +
19 files changed, 696 insertions(+), 221 deletions(-)
diff --git a/docs/generated/flink_connector_configuration.html
b/docs/generated/flink_connector_configuration.html
index 8f5b2eb9c3..3d7c5149a4 100644
--- a/docs/generated/flink_connector_configuration.html
+++ b/docs/generated/flink_connector_configuration.html
@@ -282,7 +282,7 @@ under the License.
<td><h5>sink.cross-partition.managed-memory</h5></td>
<td style="word-wrap: break-word;">256 mb</td>
<td>MemorySize</td>
- <td>Weight of managed memory for RocksDB in cross-partition
update, Flink will compute the memory size according to the weight, the actual
memory used depends on the running environment.</td>
+ <td>Weight of managed memory for the local key-value index in
cross-partition update, Flink will compute the memory size according to the
weight, the actual memory used depends on the running environment.</td>
</tr>
<tr>
<td><h5>sink.key-only-deletes.enabled</h5></td>
diff --git
a/paimon-benchmark/paimon-micro-benchmarks/src/test/java/org/apache/paimon/benchmark/lookup/LocalKvDbBenchmark.java
b/paimon-benchmark/paimon-micro-benchmarks/src/test/java/org/apache/paimon/benchmark/lookup/LocalKvDbBenchmark.java
index a7ddaa8305..b02971f60f 100644
---
a/paimon-benchmark/paimon-micro-benchmarks/src/test/java/org/apache/paimon/benchmark/lookup/LocalKvDbBenchmark.java
+++
b/paimon-benchmark/paimon-micro-benchmarks/src/test/java/org/apache/paimon/benchmark/lookup/LocalKvDbBenchmark.java
@@ -18,12 +18,20 @@
package org.apache.paimon.benchmark.lookup;
+import org.apache.paimon.CoreOptions;
import org.apache.paimon.benchmark.Benchmark;
import org.apache.paimon.compression.CompressOptions;
import org.apache.paimon.data.BinaryString;
import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.data.serializer.IntSerializer;
import org.apache.paimon.data.serializer.RowCompactedSerializer;
import org.apache.paimon.io.cache.CacheManager;
+import org.apache.paimon.lookup.BulkLoader;
+import org.apache.paimon.lookup.ListBulkLoader;
+import org.apache.paimon.lookup.ListState;
+import org.apache.paimon.lookup.SetState;
+import org.apache.paimon.lookup.StateFactory;
+import org.apache.paimon.lookup.local.LocalKvStateFactory;
import org.apache.paimon.lookup.rocksdb.RocksDBBulkLoader;
import org.apache.paimon.lookup.rocksdb.RocksDBOptions;
import org.apache.paimon.lookup.rocksdb.RocksDBStateFactory;
@@ -47,7 +55,9 @@ import java.io.File;
import java.io.IOException;
import java.nio.file.Path;
import java.util.AbstractMap;
+import java.util.ArrayList;
import java.util.Iterator;
+import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.UUID;
@@ -64,6 +74,10 @@ public class LocalKvDbBenchmark {
private static final int VALUE_SIZE =
intProperty("local-kv-db.benchmark.value-size", 64);
private static final int CACHE_SIZE_MB =
intProperty("local-kv-db.benchmark.cache-size-mb", 128);
+ private static final boolean CACHE_OFF_HEAP =
+ Boolean.parseBoolean(
+ System.getProperties()
+
.getProperty("local-kv-db.benchmark.cache-off-heap", "false"));
private static final int MEMTABLE_SIZE_MB =
intProperty("local-kv-db.benchmark.memtable-size-mb", 64);
private static final int SST_FILE_SIZE_MB =
@@ -76,6 +90,11 @@ public class LocalKvDbBenchmark {
CompressOptions.defaultOptions().compress());
private static final double BLOOM_FILTER_FPP =
doubleProperty("local-kv-db.benchmark.bloom-filter-fpp", -1);
+ private static final int STATE_KEY_COUNT =
+ intProperty("local-kv-db.benchmark.state-keys", 2_000);
+ private static final int STATE_FAN_OUT =
intProperty("local-kv-db.benchmark.state-fan-out", 64);
+ private static final int STATE_CACHE_ROWS =
+ intProperty("local-kv-db.benchmark.state-cache-rows", 10_000);
private static final RowType CLUSTERING_KEY_TYPE =
RowType.of(DataTypes.BIGINT(), DataTypes.BIGINT(),
DataTypes.STRING());
private static final BinaryString[] REGIONS = createRegions();
@@ -312,6 +331,278 @@ public class LocalKvDbBenchmark {
assertThat(checksum).isNotZero();
}
+ @Test
+ public void testStateFanOutComparison() throws IOException {
+ benchmarkListState("local-list-warmup", true);
+ benchmarkListState("rocks-list-warmup", false);
+ benchmarkListBulkLoad("local-list-bulk-warmup", true);
+ benchmarkListBulkLoad("rocks-list-bulk-warmup", false);
+ benchmarkSetState("local-set-warmup", true);
+ benchmarkSetState("rocks-set-warmup", false);
+
+ StateBenchmarkResult localList = benchmarkListState("local-list",
true);
+ StateBenchmarkResult rocksList = benchmarkListState("rocks-list",
false);
+ StateBenchmarkResult localListBulk =
benchmarkListBulkLoad("local-list-bulk", true);
+ StateBenchmarkResult rocksListBulk =
benchmarkListBulkLoad("rocks-list-bulk", false);
+ StateBenchmarkResult localSet = benchmarkSetState("local-set", true);
+ StateBenchmarkResult rocksSet = benchmarkSetState("rocks-set", false);
+
+ System.out.printf(
+ Locale.ROOT,
+ "ListState (%d keys x %d values, local cache=%s):%n"
+ + " local-kv-db add=%.1f ms, get=%.1f ms, cached=%.1f
ms, close=%.1f ms, size=%.2f MB%n"
+ + " rocks add=%.1f ms, get=%.1f ms, cached=%.1f
ms, close=%.1f ms, size=%.2f MB%n",
+ STATE_KEY_COUNT,
+ STATE_FAN_OUT,
+ CACHE_OFF_HEAP ? "off-heap" : "heap",
+ localList.updateNanos / 1_000_000.0,
+ localList.readNanos / 1_000_000.0,
+ localList.cachedReadNanos / 1_000_000.0,
+ localList.closeNanos / 1_000_000.0,
+ localList.directoryBytes / (1024.0 * 1024.0),
+ rocksList.updateNanos / 1_000_000.0,
+ rocksList.readNanos / 1_000_000.0,
+ rocksList.cachedReadNanos / 1_000_000.0,
+ rocksList.closeNanos / 1_000_000.0,
+ rocksList.directoryBytes / (1024.0 * 1024.0));
+ System.out.printf(
+ Locale.ROOT,
+ "ListState bulk load (%d keys x %d values):%n"
+ + " local-kv-db load=%.1f ms, get=%.1f ms,
cached=%.1f ms, close=%.1f ms, size=%.2f MB%n"
+ + " rocks load=%.1f ms, get=%.1f ms,
cached=%.1f ms, close=%.1f ms, size=%.2f MB%n",
+ STATE_KEY_COUNT,
+ STATE_FAN_OUT,
+ localListBulk.updateNanos / 1_000_000.0,
+ localListBulk.readNanos / 1_000_000.0,
+ localListBulk.cachedReadNanos / 1_000_000.0,
+ localListBulk.closeNanos / 1_000_000.0,
+ localListBulk.directoryBytes / (1024.0 * 1024.0),
+ rocksListBulk.updateNanos / 1_000_000.0,
+ rocksListBulk.readNanos / 1_000_000.0,
+ rocksListBulk.cachedReadNanos / 1_000_000.0,
+ rocksListBulk.closeNanos / 1_000_000.0,
+ rocksListBulk.directoryBytes / (1024.0 * 1024.0));
+ System.out.printf(
+ Locale.ROOT,
+ "SetState (%d keys x %d values):%n"
+ + " local-kv-db add=%.1f ms, get=%.1f ms, cached=%.1f
ms, close=%.1f ms, size=%.2f MB%n"
+ + " rocks add=%.1f ms, get=%.1f ms, cached=%.1f
ms, close=%.1f ms, size=%.2f MB%n",
+ STATE_KEY_COUNT,
+ STATE_FAN_OUT,
+ localSet.updateNanos / 1_000_000.0,
+ localSet.readNanos / 1_000_000.0,
+ localSet.cachedReadNanos / 1_000_000.0,
+ localSet.closeNanos / 1_000_000.0,
+ localSet.directoryBytes / (1024.0 * 1024.0),
+ rocksSet.updateNanos / 1_000_000.0,
+ rocksSet.readNanos / 1_000_000.0,
+ rocksSet.cachedReadNanos / 1_000_000.0,
+ rocksSet.closeNanos / 1_000_000.0,
+ rocksSet.directoryBytes / (1024.0 * 1024.0));
+
+ checksum +=
+ localList.checksum
+ + rocksList.checksum
+ + localListBulk.checksum
+ + rocksListBulk.checksum
+ + localSet.checksum
+ + rocksSet.checksum;
+ assertThat(checksum).isNotZero();
+ }
+
+ private StateBenchmarkResult benchmarkListState(String name, boolean
local) throws IOException {
+ File directory = new File(tempDir.toFile(), name);
+ StateFactory factory = createStateFactory(directory, local);
+ long updateNanos = 0;
+ long readNanos = 0;
+ long cachedReadNanos = 0;
+ long closeNanos = 0;
+ long directoryBytes = 0;
+ long localChecksum = 0;
+ try {
+ ListState<Integer, Integer> state =
+ factory.listState(
+ "list",
+ IntSerializer.INSTANCE,
+ IntSerializer.INSTANCE,
+ STATE_CACHE_ROWS);
+ long updateStart = System.nanoTime();
+ for (int key = 0; key < STATE_KEY_COUNT; key++) {
+ for (int value = 0; value < STATE_FAN_OUT; value++) {
+ state.add(key, value);
+ }
+ }
+ updateNanos = System.nanoTime() - updateStart;
+
+ long readStart = System.nanoTime();
+ for (int key = 0; key < STATE_KEY_COUNT; key++) {
+ List<Integer> values = state.get(key);
+ if (values.size() != STATE_FAN_OUT) {
+ throw new IllegalStateException("Unexpected ListState
fan-out.");
+ }
+ localChecksum += values.size() + values.get(values.size() - 1);
+ }
+ readNanos = System.nanoTime() - readStart;
+
+ long cachedReadStart = System.nanoTime();
+ for (int key = 0; key < STATE_KEY_COUNT; key++) {
+ List<Integer> values = state.get(key);
+ if (values.size() != STATE_FAN_OUT) {
+ throw new IllegalStateException("Unexpected cached
ListState fan-out.");
+ }
+ localChecksum += values.size() + values.get(values.size() - 1);
+ }
+ cachedReadNanos = System.nanoTime() - cachedReadStart;
+ } finally {
+ long closeStart = System.nanoTime();
+ try {
+ factory.close();
+ } finally {
+ closeNanos = System.nanoTime() - closeStart;
+ directoryBytes = directorySize(directory);
+ FileIOUtils.deleteDirectoryQuietly(directory);
+ }
+ }
+ return new StateBenchmarkResult(
+ updateNanos, readNanos, cachedReadNanos, closeNanos,
directoryBytes, localChecksum);
+ }
+
+ private StateBenchmarkResult benchmarkSetState(String name, boolean local)
throws IOException {
+ File directory = new File(tempDir.toFile(), name);
+ StateFactory factory = createStateFactory(directory, local);
+ long updateNanos = 0;
+ long readNanos = 0;
+ long cachedReadNanos = 0;
+ long closeNanos = 0;
+ long directoryBytes = 0;
+ long localChecksum = 0;
+ try {
+ SetState<Integer, Integer> state =
+ factory.setState(
+ "set",
+ IntSerializer.INSTANCE,
+ IntSerializer.INSTANCE,
+ STATE_CACHE_ROWS);
+ long updateStart = System.nanoTime();
+ for (int key = 0; key < STATE_KEY_COUNT; key++) {
+ for (int value = 0; value < STATE_FAN_OUT; value++) {
+ state.add(key, value);
+ }
+ }
+ updateNanos = System.nanoTime() - updateStart;
+
+ long readStart = System.nanoTime();
+ for (int key = 0; key < STATE_KEY_COUNT; key++) {
+ List<Integer> values = state.get(key);
+ if (values.size() != STATE_FAN_OUT) {
+ throw new IllegalStateException("Unexpected SetState
fan-out.");
+ }
+ localChecksum += values.size() + values.get(values.size() - 1);
+ }
+ readNanos = System.nanoTime() - readStart;
+
+ long cachedReadStart = System.nanoTime();
+ for (int key = 0; key < STATE_KEY_COUNT; key++) {
+ List<Integer> values = state.get(key);
+ if (values.size() != STATE_FAN_OUT) {
+ throw new IllegalStateException("Unexpected cached
SetState fan-out.");
+ }
+ localChecksum += values.size() + values.get(values.size() - 1);
+ }
+ cachedReadNanos = System.nanoTime() - cachedReadStart;
+ } finally {
+ long closeStart = System.nanoTime();
+ try {
+ factory.close();
+ } finally {
+ closeNanos = System.nanoTime() - closeStart;
+ directoryBytes = directorySize(directory);
+ FileIOUtils.deleteDirectoryQuietly(directory);
+ }
+ }
+ return new StateBenchmarkResult(
+ updateNanos, readNanos, cachedReadNanos, closeNanos,
directoryBytes, localChecksum);
+ }
+
+ private StateBenchmarkResult benchmarkListBulkLoad(String name, boolean
local)
+ throws IOException {
+ File directory = new File(tempDir.toFile(), name);
+ StateFactory factory = createStateFactory(directory, local);
+ long loadNanos = 0;
+ long readNanos = 0;
+ long cachedReadNanos = 0;
+ long closeNanos = 0;
+ long directoryBytes = 0;
+ long localChecksum = 0;
+ try {
+ ListState<Integer, Integer> state =
+ factory.listState(
+ "list",
+ IntSerializer.INSTANCE,
+ IntSerializer.INSTANCE,
+ STATE_CACHE_ROWS);
+ ListBulkLoader loader = state.createBulkLoader();
+ long loadStart = System.nanoTime();
+ for (int key = 0; key < STATE_KEY_COUNT; key++) {
+ List<byte[]> values = new ArrayList<>(STATE_FAN_OUT);
+ for (int value = 0; value < STATE_FAN_OUT; value++) {
+ values.add(state.serializeValue(value));
+ }
+ loader.write(state.serializeKey(key), values);
+ }
+ loader.finish();
+ loadNanos = System.nanoTime() - loadStart;
+
+ long readStart = System.nanoTime();
+ for (int key = 0; key < STATE_KEY_COUNT; key++) {
+ List<Integer> values = state.get(key);
+ if (values.size() != STATE_FAN_OUT) {
+ throw new IllegalStateException("Unexpected bulk-loaded
ListState fan-out.");
+ }
+ localChecksum += values.size() + values.get(values.size() - 1);
+ }
+ readNanos = System.nanoTime() - readStart;
+
+ long cachedReadStart = System.nanoTime();
+ for (int key = 0; key < STATE_KEY_COUNT; key++) {
+ List<Integer> values = state.get(key);
+ if (values.size() != STATE_FAN_OUT) {
+ throw new IllegalStateException(
+ "Unexpected cached bulk-loaded ListState
fan-out.");
+ }
+ localChecksum += values.size() + values.get(values.size() - 1);
+ }
+ cachedReadNanos = System.nanoTime() - cachedReadStart;
+ } catch (BulkLoader.WriteException e) {
+ throw new IOException(e);
+ } finally {
+ long closeStart = System.nanoTime();
+ try {
+ factory.close();
+ } finally {
+ closeNanos = System.nanoTime() - closeStart;
+ directoryBytes = directorySize(directory);
+ FileIOUtils.deleteDirectoryQuietly(directory);
+ }
+ }
+ return new StateBenchmarkResult(
+ loadNanos, readNanos, cachedReadNanos, closeNanos,
directoryBytes, localChecksum);
+ }
+
+ private StateFactory createStateFactory(File directory, boolean local)
throws IOException {
+ Options options = new Options();
+ options.set(
+ CoreOptions.LOOKUP_CACHE_MAX_MEMORY_SIZE,
MemorySize.ofMebiBytes(CACHE_SIZE_MB));
+ options.set(CoreOptions.LOOKUP_CACHE_SPILL_COMPRESSION, COMPRESSION);
+ options.set(RocksDBOptions.BLOCK_CACHE_SIZE,
MemorySize.ofMebiBytes(CACHE_SIZE_MB));
+ options.set(RocksDBOptions.COMPRESSION_TYPE, rocksCompression());
+ if (local) {
+ return new LocalKvStateFactory(
+ directory.getAbsolutePath(), options, null, null,
CACHE_OFF_HEAP);
+ }
+ return new RocksDBStateFactory(directory.getAbsolutePath(), options,
null);
+ }
+
private void mixedReadWrite(LocalKvDb db, byte[][] keys, byte[][] values) {
try {
long localChecksum = 0;
@@ -408,7 +699,7 @@ public class LocalKvDbBenchmark {
.memTableFlushThreshold(MEMTABLE_SIZE_MB * 1024L *
1024L)
.maxSstFileSize(SST_FILE_SIZE_MB * 1024L * 1024L)
.blockSize(BLOCK_SIZE_KB * 1024)
- .cacheManager(new
CacheManager(MemorySize.ofMebiBytes(CACHE_SIZE_MB), 0))
+ .cacheManager(createCacheManager(CACHE_SIZE_MB))
.compressOptions(new CompressOptions(COMPRESSION, 1))
.bloomFilterEnabled(BLOOM_FILTER_FPP > 0);
if (BLOOM_FILTER_FPP > 0) {
@@ -422,7 +713,7 @@ public class LocalKvDbBenchmark {
.memTableFlushThreshold(MEMTABLE_SIZE_MB * 1024L * 1024L)
.maxSstFileSize(SST_FILE_SIZE_MB * 1024L * 1024L)
.blockSize(blockSizeKb * 1024)
- .cacheManager(new
CacheManager(MemorySize.ofMebiBytes(cacheSizeMb), 0))
+ .cacheManager(createCacheManager(cacheSizeMb))
.compressOptions(CompressOptions.defaultOptions())
.bloomFilterEnabled(true)
.bloomFilterFpp(0.1)
@@ -431,6 +722,13 @@ public class LocalKvDbBenchmark {
.build();
}
+ private CacheManager createCacheManager(int cacheSizeMb) {
+ MemorySize cacheSize = MemorySize.ofMebiBytes(cacheSizeMb);
+ return CACHE_OFF_HEAP
+ ? CacheManager.createOffHeap(cacheSize, 0)
+ : new CacheManager(cacheSize, 0);
+ }
+
private RocksDbHandle createRocksDb(File directory) throws IOException {
Options options = new Options();
options.set(RocksDBOptions.WRITE_BUFFER_SIZE,
MemorySize.ofMebiBytes(MEMTABLE_SIZE_MB));
@@ -602,10 +900,11 @@ public class LocalKvDbBenchmark {
private static String benchmarkDescription() {
return String.format(
Locale.ROOT,
-
"%d-records-%dB-value-%dMB-caffeine-cache-%dMB-memtable-%dMB-sst-%dKB-block-%s-bloom-%s",
+
"%d-records-%dB-value-%dMB-%s-caffeine-cache-%dMB-memtable-%dMB-sst-%dKB-block-%s-bloom-%s",
RECORD_COUNT,
VALUE_SIZE,
CACHE_SIZE_MB,
+ CACHE_OFF_HEAP ? "off-heap" : "heap",
MEMTABLE_SIZE_MB,
SST_FILE_SIZE_MB,
BLOCK_SIZE_KB,
@@ -646,4 +945,29 @@ public class LocalKvDbBenchmark {
factory.close();
}
}
+
+ private static class StateBenchmarkResult {
+
+ private final long updateNanos;
+ private final long readNanos;
+ private final long cachedReadNanos;
+ private final long closeNanos;
+ private final long directoryBytes;
+ private final long checksum;
+
+ private StateBenchmarkResult(
+ long updateNanos,
+ long readNanos,
+ long cachedReadNanos,
+ long closeNanos,
+ long directoryBytes,
+ long checksum) {
+ this.updateNanos = updateNanos;
+ this.readNanos = readNanos;
+ this.cachedReadNanos = cachedReadNanos;
+ this.closeNanos = closeNanos;
+ this.directoryBytes = directoryBytes;
+ this.checksum = checksum;
+ }
+ }
}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/lookup/sort/db/LocalKvDb.java
b/paimon-common/src/main/java/org/apache/paimon/lookup/sort/db/LocalKvDb.java
index 3b976a67b5..bf2e5698f7 100644
---
a/paimon-common/src/main/java/org/apache/paimon/lookup/sort/db/LocalKvDb.java
+++
b/paimon-common/src/main/java/org/apache/paimon/lookup/sort/db/LocalKvDb.java
@@ -43,8 +43,8 @@ import java.util.AbstractMap;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Comparator;
-import java.util.HashMap;
import java.util.Iterator;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.PriorityQueue;
@@ -95,6 +95,19 @@ public class LocalKvDb implements Closeable {
/** Marker used when an SST's entry count cannot be estimated before
writing. */
static final long UNKNOWN_NUM_ENTRIES = -1;
+ /** Bound open SST inputs while retaining one hot reader per LSM level. */
+ static final int MAX_CACHED_READERS = MAX_LEVELS;
+
+ /** Initial number of values retained by a lazy MemTable merge. */
+ private static final int MERGE_VALUES_INITIAL_CAPACITY = 32;
+
+ /** Approximate retained memory for a lazy merge function and its initial
reference array. */
+ private static final long MEMTABLE_MERGE_FUNCTION_OVERHEAD =
+ 64L + MERGE_VALUES_INITIAL_CAPACITY * Long.BYTES;
+
+ /** Approximate array header and reference overhead for each retained
merge value. */
+ private static final long MEMTABLE_MERGED_VALUE_OVERHEAD = 24;
+
/**
* Estimated per-entry memory overhead in the MemTable's TreeMap, beyond
the raw key/value
* bytes. This accounts for:
@@ -121,8 +134,8 @@ public class LocalKvDb implements Closeable {
private final LsmCompactor compaction;
@Nullable private final MergeOperator mergeOperator;
- /** Active MemTable: key -> value bytes (empty byte[] = tombstone). */
- private TreeMap<MemorySlice, byte[]> memTable;
+ /** Active MemTable: key -> value bytes, tombstone, or lazy merge
function. */
+ private TreeMap<MemorySlice, Object> memTable;
/** Estimated size of the current MemTable in bytes. */
private long memTableSize;
@@ -157,7 +170,7 @@ public class LocalKvDb implements Closeable {
this.memTable = new TreeMap<>(keyComparator);
this.memTableSize = 0;
this.levels = new LsmLevels(MAX_LEVELS);
- this.readerCache = new HashMap<>();
+ this.readerCache = new LinkedHashMap<>(16, 0.75f, true);
this.fileSequence = new AtomicLong();
this.activeBulkLoadWriter = null;
this.openRangeIterators = 0;
@@ -240,10 +253,31 @@ public class LocalKvDb implements Closeable {
+ "Use delete() to remove a key.");
}
MemorySlice wrappedKey = MemorySlice.wrap(key);
- byte[] oldValue = memTable.put(wrappedKey, value);
+ if (mergeOperator != null) {
+ Map.Entry<MemorySlice, Object> previous =
memTable.lowerEntry(wrappedKey);
+ if (previous != null
+ && !isMemTableTombstone(previous.getValue())
+ && mergeOperator.canMerge(previous.getKey(), wrappedKey)) {
+ Object previousValue = previous.getValue();
+ long previousSize = estimatedMemTableValueSize(previousValue);
+ MemTableMergeFunction mergeFunction;
+ if (previousValue instanceof MemTableMergeFunction) {
+ mergeFunction = (MemTableMergeFunction) previousValue;
+ } else {
+ mergeFunction = new MemTableMergeFunction(mergeOperator);
+ mergeFunction.reset((byte[]) previousValue);
+ }
+ mergeFunction.add(value);
+ memTableSize += mergeFunction.estimatedSize() - previousSize;
+ memTable.put(previous.getKey(), mergeFunction);
+ maybeFlushMemTable();
+ return;
+ }
+ }
+ Object oldValue = memTable.put(wrappedKey, value);
long delta = key.length + value.length;
if (oldValue != null) {
- delta -= (key.length + oldValue.length);
+ delta -= (key.length + estimatedMemTableValueSize(oldValue));
} else {
delta += PER_ENTRY_OVERHEAD;
}
@@ -262,10 +296,10 @@ public class LocalKvDb implements Closeable {
ensureNoRangeIterator();
checkCompactionFailure();
MemorySlice wrappedKey = MemorySlice.wrap(key);
- byte[] oldValue = memTable.put(wrappedKey, TOMBSTONE);
+ Object oldValue = memTable.put(wrappedKey, TOMBSTONE);
long delta = key.length;
if (oldValue != null) {
- delta -= (key.length + oldValue.length);
+ delta -= (key.length + estimatedMemTableValueSize(oldValue));
} else {
delta += PER_ENTRY_OVERHEAD;
}
@@ -556,7 +590,8 @@ public class LocalKvDb implements Closeable {
// 1. Search MemTable first (newest data)
MemorySlice wrappedKey = MemorySlice.wrap(key);
- byte[] memValue = memTable.get(wrappedKey);
+ Object memTableValue = memTable.get(wrappedKey);
+ byte[] memValue = memTableValue == null ? null :
materializeMemTableValue(memTableValue);
if (memValue != null) {
return isTombstone(memValue) ? null : memValue;
}
@@ -624,7 +659,7 @@ public class LocalKvDb implements Closeable {
to == null || keyComparator.compare(from, to) <= 0,
"Range start must not be greater than range end.");
- Map<MemorySlice, byte[]> memoryEntries =
+ Map<MemorySlice, Object> memoryEntries =
to == null ? memTable.tailMap(from, true) :
memTable.subMap(from, true, to, false);
LsmLevels.RangeSnapshot snapshot = levels.openRangeSnapshot(from, to,
keyComparator);
try {
@@ -663,7 +698,7 @@ public class LocalKvDb implements Closeable {
}
private void flushMemTable() throws IOException {
- TreeMap<MemorySlice, byte[]> snapshot = memTable;
+ TreeMap<MemorySlice, Object> snapshot = memTable;
SstFileMetadata metadata = writeMemTableToSst(snapshot);
levels.addLevelZeroFile(metadata);
memTable = new TreeMap<>(keyComparator);
@@ -800,17 +835,57 @@ public class LocalKvDb implements Closeable {
}
}
+ private boolean isMemTableTombstone(Object value) {
+ return value instanceof byte[] && isTombstone((byte[]) value);
+ }
+
+ private long estimatedMemTableValueSize(Object value) {
+ return value instanceof byte[]
+ ? ((byte[]) value).length
+ : ((MemTableMergeFunction) value).estimatedSize();
+ }
+
+ private byte[] materializeMemTableValue(Object value) throws IOException {
+ if (value instanceof byte[]) {
+ return (byte[]) value;
+ }
+
+ MemTableMergeFunction mergeFunction = (MemTableMergeFunction) value;
+ long previousSize = mergeFunction.estimatedSize();
+ byte[] result = mergeFunction.getResult();
+ memTableSize += mergeFunction.estimatedSize() - previousSize;
+ return result;
+ }
+
@Nullable
private byte[] lookupInFile(File file, byte[] key) throws IOException {
+ return getOrCreateReader(file).lookup(key);
+ }
+
+ private SortLookupStoreReader getOrCreateReader(File file) throws
IOException {
SortLookupStoreReader reader = readerCache.get(file);
if (reader == null) {
reader = storeFactory.createReader(file);
+ if (readerCache.size() >= MAX_CACHED_READERS) {
+ Iterator<Map.Entry<File, SortLookupStoreReader>> iterator =
+ readerCache.entrySet().iterator();
+ Map.Entry<File, SortLookupStoreReader> eldest =
iterator.next();
+ iterator.remove();
+ try {
+ eldest.getValue().closeInput();
+ } catch (IOException e) {
+ LOG.warn(
+ "Failed to close evicted reader for SST file: {}",
+ eldest.getKey().getName(),
+ e);
+ }
+ }
readerCache.put(file, reader);
}
- return reader.lookup(key);
+ return reader;
}
- private SstFileMetadata writeMemTableToSst(TreeMap<MemorySlice, byte[]>
data)
+ private SstFileMetadata writeMemTableToSst(TreeMap<MemorySlice, Object>
data)
throws IOException {
File sstFile = newSstFile();
SortLookupStoreWriter writer = null;
@@ -823,8 +898,8 @@ public class LocalKvDb implements Closeable {
SstMetadataWriter output = new SstMetadataWriter(writer);
RecordCombiningWriter combiningWriter =
new RecordCombiningWriter(mergeOperator, output);
- for (Map.Entry<MemorySlice, byte[]> entry : data.entrySet()) {
- combiningWriter.put(entry.getKey(), entry.getValue());
+ for (Map.Entry<MemorySlice, Object> entry : data.entrySet()) {
+ combiningWriter.put(entry.getKey(),
materializeMemTableValue(entry.getValue()));
}
combiningWriter.finish();
writer.close();
@@ -910,7 +985,7 @@ public class LocalKvDb implements Closeable {
private RangeIterator(
LsmLevels.RangeSnapshot snapshot,
- Map<MemorySlice, byte[]> memoryEntries,
+ Map<MemorySlice, Object> memoryEntries,
byte[] fromInclusive,
@Nullable MemorySlice toExclusive)
throws IOException {
@@ -1011,14 +1086,15 @@ public class LocalKvDb implements Closeable {
MemorySlice value();
}
- private static final class MemoryRangeSource implements RangeSource {
+ private final class MemoryRangeSource implements RangeSource {
private final int priority;
- private final Iterator<Map.Entry<MemorySlice, byte[]>> iterator;
+ private final Iterator<Map.Entry<MemorySlice, Object>> iterator;
- @Nullable private Map.Entry<MemorySlice, byte[]> current;
+ @Nullable private Map.Entry<MemorySlice, Object> current;
+ @Nullable private byte[] currentValue;
- private MemoryRangeSource(int priority,
Iterator<Map.Entry<MemorySlice, byte[]>> iterator) {
+ private MemoryRangeSource(int priority,
Iterator<Map.Entry<MemorySlice, Object>> iterator) {
this.priority = priority;
this.iterator = iterator;
}
@@ -1029,8 +1105,9 @@ public class LocalKvDb implements Closeable {
}
@Override
- public boolean advance() {
+ public boolean advance() throws IOException {
current = iterator.hasNext() ? iterator.next() : null;
+ currentValue = current == null ? null :
materializeMemTableValue(current.getValue());
return current != null;
}
@@ -1041,7 +1118,7 @@ public class LocalKvDb implements Closeable {
@Override
public MemorySlice value() {
- return MemorySlice.wrap(current.getValue());
+ return MemorySlice.wrap(currentValue);
}
}
@@ -1099,29 +1176,26 @@ public class LocalKvDb implements Closeable {
return null;
}
- SortLookupStoreReader reader = storeFactory.createReader(file);
- try (Closeable ignored = reader::closeInput) {
- SstFileReader.SstFileIterator iterator =
reader.createIterator();
- byte[] seekKey =
- resumeAfterKey == null ? fromInclusive :
resumeAfterKey.copyBytes();
- iterator.seekTo(seekKey);
- boolean skipResumeKey = resumeAfterKey != null;
- while (true) {
- BlockIterator nextBlock = iterator.readBatch();
- if (nextBlock == null) {
- finished = true;
- return null;
- }
+ SortLookupStoreReader reader = getOrCreateReader(file);
+ SstFileReader.SstFileIterator iterator = reader.createIterator();
+ byte[] seekKey = resumeAfterKey == null ? fromInclusive :
resumeAfterKey.copyBytes();
+ iterator.seekTo(seekKey);
+ boolean skipResumeKey = resumeAfterKey != null;
+ while (true) {
+ BlockIterator nextBlock = iterator.readBatch();
+ if (nextBlock == null) {
+ finished = true;
+ return null;
+ }
- if (skipResumeKey) {
- if (nextBlock.seekTo(resumeAfterKey)) {
- nextBlock.next();
- }
- skipResumeKey = false;
- }
- if (nextBlock.hasNext()) {
- return nextBlock;
+ if (skipResumeKey) {
+ if (nextBlock.seekTo(resumeAfterKey)) {
+ nextBlock.next();
}
+ skipResumeKey = false;
+ }
+ if (nextBlock.hasNext()) {
+ return nextBlock;
}
}
}
@@ -1141,6 +1215,60 @@ public class LocalKvDb implements Closeable {
return value.length() == 0;
}
+ /**
+ * Retains raw merge operands until a read or flush requires their
serialized result.
+ *
+ * <p>After materialization, the operands are reset to the result so
repeated reads do not merge
+ * again and later appends only retain the previous result plus new
operands.
+ */
+ private static final class MemTableMergeFunction {
+
+ private final MergeOperator mergeOperator;
+ private List<byte[]> values;
+
+ @Nullable private byte[] result;
+ private long retainedValueBytes;
+
+ private MemTableMergeFunction(MergeOperator mergeOperator) {
+ this.mergeOperator = mergeOperator;
+ this.values = new ArrayList<>(MERGE_VALUES_INITIAL_CAPACITY);
+ }
+
+ private void reset(byte[] value) {
+ if (values.size() > MERGE_VALUES_INITIAL_CAPACITY) {
+ values = new ArrayList<>(MERGE_VALUES_INITIAL_CAPACITY);
+ } else {
+ values.clear();
+ }
+ values.add(value);
+ result = value;
+ retainedValueBytes = value.length;
+ }
+
+ private void add(byte[] value) {
+ values.add(value);
+ result = null;
+ retainedValueBytes += value.length;
+ }
+
+ private byte[] getResult() throws IOException {
+ if (result == null) {
+ byte[] merged = mergeOperator.merge(values);
+ if (isTombstone(merged)) {
+ throw new IllegalStateException("MergeOperator returned
the tombstone marker.");
+ }
+ reset(merged);
+ }
+ return result;
+ }
+
+ private long estimatedSize() {
+ return MEMTABLE_MERGE_FUNCTION_OVERHEAD
+ + retainedValueBytes
+ + values.size() * MEMTABLE_MERGED_VALUE_OVERHEAD;
+ }
+ }
+
/** Callback for visiting an entry during a range scan. */
@FunctionalInterface
public interface RangeEntryConsumer {
@@ -1151,24 +1279,15 @@ public class LocalKvDb implements Closeable {
/**
* Operator for combining adjacent logical records while flushing and
compacting SST files.
*
- * <p>MemTable writes remain independent so merge-heavy workloads do not
pay repeated
- * read-modify-write costs. The first record's key is retained for the
combined value.
+ * <p>The first record's key is retained for the combined value.
+ *
+ * <p>Mergeable physical keys must be append-only and must not be reused
after being consumed
+ * into the first record.
*/
public interface MergeOperator {
boolean canMerge(MemorySlice firstKey, MemorySlice nextKey);
- /**
- * Return whether a tombstone can be absorbed into the pending merge
group.
- *
- * <p>Tombstones are merge boundaries by default. Operators which
produce tombstones for
- * consumed physical keys can opt in so later compactions can merge
across those synthetic
- * tombstones.
- */
- default boolean canMergeTombstone(MemorySlice firstKey, MemorySlice
tombstoneKey) {
- return false;
- }
-
byte[] merge(List<byte[]> values) throws IOException;
}
diff --git
a/paimon-common/src/main/java/org/apache/paimon/lookup/sort/db/RecordCombiningWriter.java
b/paimon-common/src/main/java/org/apache/paimon/lookup/sort/db/RecordCombiningWriter.java
index 224ef16aa8..edf5eedd0d 100644
---
a/paimon-common/src/main/java/org/apache/paimon/lookup/sort/db/RecordCombiningWriter.java
+++
b/paimon-common/src/main/java/org/apache/paimon/lookup/sort/db/RecordCombiningWriter.java
@@ -55,12 +55,8 @@ final class RecordCombiningWriter {
}
if (isTombstone(value)) {
- if (pendingKey != null &&
mergeOperator.canMergeTombstone(pendingKey, key)) {
- pendingKeys.add(MemorySlice.wrap(key.copyBytes()));
- } else {
- flushPending();
- consumer.accept(key, value);
- }
+ flushPending();
+ consumer.accept(key, value);
return;
}
diff --git
a/paimon-common/src/test/java/org/apache/paimon/lookup/sort/db/LocalKvDbTest.java
b/paimon-common/src/test/java/org/apache/paimon/lookup/sort/db/LocalKvDbTest.java
index d2673e5be2..9f9bc8e4bd 100644
---
a/paimon-common/src/test/java/org/apache/paimon/lookup/sort/db/LocalKvDbTest.java
+++
b/paimon-common/src/test/java/org/apache/paimon/lookup/sort/db/LocalKvDbTest.java
@@ -40,6 +40,8 @@ import java.util.Collections;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
import static java.nio.charset.StandardCharsets.UTF_8;
@@ -277,7 +279,7 @@ public class LocalKvDbTest {
}
@Test
- public void testRangeIteratorDoesNotCacheReadersForOverlappingSsts()
throws IOException {
+ public void testRangeIteratorBoundsCachedReadersForOverlappingSsts()
throws IOException {
int fileCount = 128;
File directory = new File(tempDir.toFile(),
"range-iterator-many-ssts");
try (LocalKvDb db =
@@ -294,14 +296,14 @@ public class LocalKvDbTest {
int entryCount = 0;
try (LocalKvDb.RangeIterator iterator =
db.rangeIterator("key-00000".getBytes(UTF_8),
"key-99999".getBytes(UTF_8))) {
- Assertions.assertEquals(0, db.getCachedReaderCount());
+ Assertions.assertEquals(LocalKvDb.MAX_CACHED_READERS,
db.getCachedReaderCount());
while (iterator.advanceNext()) {
entryCount++;
}
}
Assertions.assertEquals(fileCount, entryCount);
- Assertions.assertEquals(0, db.getCachedReaderCount());
+ Assertions.assertEquals(LocalKvDb.MAX_CACHED_READERS,
db.getCachedReaderCount());
}
}
@@ -563,54 +565,6 @@ public class LocalKvDbTest {
}
}
- @Test
- public void testCompactionMergesAcrossAbsorbedTombstones() throws
IOException {
- File directory = new File(tempDir.toFile(), "tombstone-merge-db");
- LocalKvDb.MergeOperator mergeOperator =
- new LocalKvDb.MergeOperator() {
- @Override
- public boolean canMerge(MemorySlice firstKey, MemorySlice
nextKey) {
- return firstKey.readByte(0) == nextKey.readByte(0);
- }
-
- @Override
- public boolean canMergeTombstone(
- MemorySlice firstKey, MemorySlice tombstoneKey) {
- return canMerge(firstKey, tombstoneKey);
- }
-
- @Override
- public byte[] merge(List<byte[]> values) {
- StringBuilder merged = new StringBuilder();
- for (byte[] value : values) {
- if (merged.length() > 0) {
- merged.append('+');
- }
- merged.append(new String(value, UTF_8));
- }
- return merged.toString().getBytes(UTF_8);
- }
- };
- try (LocalKvDb db =
- LocalKvDb.builder(directory)
- .level0FileNumCompactTrigger(100)
- .compressOptions(new CompressOptions("none", 1))
- .mergeOperator(mergeOperator)
- .build()) {
- putString(db, "a-0", "one");
- putString(db, "a-1", "two");
- db.flush();
- putString(db, "a-2", "three");
- db.flush();
-
- db.compact();
-
- Assertions.assertEquals("one+two+three", getString(db, "a-0"));
- Assertions.assertNull(getString(db, "a-1"));
- Assertions.assertNull(getString(db, "a-2"));
- }
- }
-
@Test
public void
testCompactionMergesAcrossFileGroupsBeforeFilteringExpiration() throws
IOException {
File directory = new File(tempDir.toFile(),
"cross-group-expiration-merge-db");
@@ -667,8 +621,9 @@ public class LocalKvDbTest {
}
@Test
- public void testFlushMergeShadowsConsumedKeysInOlderRuns() throws
IOException {
- File directory = new File(tempDir.toFile(), "flush-merge-shadow-db");
+ public void testMemTableMergeMaterializesLazily() throws IOException {
+ AtomicBoolean failMerge = new AtomicBoolean();
+ AtomicInteger mergeCount = new AtomicInteger();
LocalKvDb.MergeOperator mergeOperator =
new LocalKvDb.MergeOperator() {
@Override
@@ -677,32 +632,51 @@ public class LocalKvDbTest {
}
@Override
- public byte[] merge(List<byte[]> values) {
- return (new String(values.get(0), UTF_8)
- + "+"
- + new String(values.get(1), UTF_8))
- .getBytes(UTF_8);
+ public byte[] merge(List<byte[]> values) throws
IOException {
+ mergeCount.incrementAndGet();
+ if (failMerge.get()) {
+ throw new IOException("Expected merge failure.");
+ }
+ StringBuilder result = new StringBuilder();
+ for (byte[] value : values) {
+ result.append(new String(value, UTF_8));
+ }
+ return result.toString().getBytes(UTF_8);
}
};
try (LocalKvDb db =
- LocalKvDb.builder(directory)
+ LocalKvDb.builder(new File(tempDir.toFile(),
"lazy-memtable-merge"))
.level0FileNumCompactTrigger(100)
.compressOptions(new CompressOptions("none", 1))
.mergeOperator(mergeOperator)
.build()) {
- List<Map.Entry<byte[], byte[]>> oldValues = new ArrayList<>();
- oldValues.add(entry("a-1", "old-1"));
- oldValues.add(entry("a-2", "old-2"));
- db.bulkLoad(oldValues.iterator(), oldValues.size());
+ putString(db, "a-0", "one");
+ putString(db, "a-1", "two");
+ putString(db, "a-2", "three");
+ Assertions.assertEquals(0, mergeCount.get());
- putString(db, "a-1", "new-1");
- putString(db, "a-2", "new-2");
- Assertions.assertEquals("new-2", getString(db, "a-2"));
+ Assertions.assertEquals("onetwothree", getString(db, "a-0"));
+ Assertions.assertEquals(1, mergeCount.get());
+ Assertions.assertEquals("onetwothree", getString(db, "a-0"));
+ Assertions.assertEquals(1, mergeCount.get());
- db.flush();
+ putString(db, "a-3", "four");
+ Assertions.assertEquals(1, mergeCount.get());
+ Assertions.assertEquals("onetwothreefour", getString(db, "a-0"));
+ Assertions.assertEquals(2, mergeCount.get());
+
+ putString(db, "a-4", "five");
+ failMerge.set(true);
+ Assertions.assertThrows(IOException.class, db::flush);
+ Assertions.assertEquals(3, mergeCount.get());
+ failMerge.set(false);
+ Assertions.assertEquals("onetwothreefourfive", getString(db,
"a-0"));
+ Assertions.assertEquals(4, mergeCount.get());
- Assertions.assertEquals("new-1+new-2", getString(db, "a-1"));
- Assertions.assertNull(getString(db, "a-2"));
+ db.flush();
+ Assertions.assertEquals(4, mergeCount.get());
+ Assertions.assertEquals("onetwothreefourfive", getString(db,
"a-0"));
+ Assertions.assertNull(getString(db, "a-1"));
}
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java
b/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java
index cb7d45a648..6e5d4b6e51 100644
---
a/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java
+++
b/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java
@@ -28,11 +28,13 @@ import
org.apache.paimon.data.serializer.InternalRowSerializer;
import org.apache.paimon.data.serializer.RowCompactedSerializer;
import org.apache.paimon.disk.IOManager;
import org.apache.paimon.disk.RowBuffer;
-import org.apache.paimon.lookup.rocksdb.RocksDBBulkLoader;
+import org.apache.paimon.lookup.BulkLoader;
+import org.apache.paimon.lookup.StateFactory;
+import org.apache.paimon.lookup.StateUtils;
+import org.apache.paimon.lookup.ValueBulkLoader;
+import org.apache.paimon.lookup.ValueState;
+import org.apache.paimon.lookup.local.LocalKvStateFactory;
import org.apache.paimon.lookup.rocksdb.RocksDBOptions;
-import org.apache.paimon.lookup.rocksdb.RocksDBState;
-import org.apache.paimon.lookup.rocksdb.RocksDBStateFactory;
-import org.apache.paimon.lookup.rocksdb.RocksDBValueState;
import org.apache.paimon.memory.HeapMemorySegmentPool;
import org.apache.paimon.options.MemorySize;
import org.apache.paimon.options.Options;
@@ -67,6 +69,7 @@ import java.util.HashMap;
import java.util.List;
import java.util.NoSuchElementException;
import java.util.UUID;
+import java.util.concurrent.ExecutorService;
import java.util.function.BiConsumer;
import java.util.function.Function;
import java.util.stream.IntStream;
@@ -97,8 +100,8 @@ public class GlobalIndexAssigner implements Serializable,
Closeable {
private transient PartitionKeyExtractor<InternalRow> extractor;
private transient PartitionKeyExtractor<InternalRow> keyPartExtractor;
private transient File path;
- private transient RocksDBStateFactory stateFactory;
- private transient RocksDBValueState<InternalRow, PositiveIntInt> keyIndex;
+ private transient StateFactory stateFactory;
+ private transient ValueState<InternalRow, PositiveIntInt> keyIndex;
private transient IDMapping<BinaryRow> partMapping;
private transient BucketAssigner bucketAssigner;
@@ -112,6 +115,7 @@ public class GlobalIndexAssigner implements Serializable,
Closeable {
public void open(
long offHeapMemory,
+ @Nullable ExecutorService compactionExecutor,
IOManager ioManager,
int numAssigners,
int assignId,
@@ -133,24 +137,30 @@ public class GlobalIndexAssigner implements Serializable,
Closeable {
this.keyPartExtractor = new
KeyPartPartitionKeyExtractor(table.schema());
String tmpDir = ioManager.pickTempDir();
- this.path = new File(tmpDir, "rocksdb-" + UUID.randomUUID());
+ this.path = new File(tmpDir, "local-kv-" + UUID.randomUUID());
if (!this.path.mkdirs()) {
throw new RuntimeException(
- "Failed to create RocksDB cache directory in temp dirs: "
+ "Failed to create local KV cache directory in temp dirs: "
+ Arrays.toString(ioManager.tempDirs()));
}
// state
Options options = coreOptions.toConfiguration();
- Options rocksdbOptions = Options.fromMap(new
HashMap<>(options.toMap()));
+ Options stateOptions = Options.fromMap(new HashMap<>(options.toMap()));
// we should avoid too small memory
- long blockCache = Math.max(offHeapMemory,
rocksdbOptions.get(BLOCK_CACHE_SIZE).getBytes());
- rocksdbOptions.set(BLOCK_CACHE_SIZE, new MemorySize(blockCache));
+ long configuredCache =
+ options.contains(CoreOptions.LOOKUP_CACHE_MAX_MEMORY_SIZE)
+ ? coreOptions.lookupCacheMaxMemory().getBytes()
+ : options.get(BLOCK_CACHE_SIZE).getBytes();
+ long cacheMemory = Math.max(offHeapMemory, configuredCache);
+ stateOptions.set(CoreOptions.LOOKUP_CACHE_MAX_MEMORY_SIZE, new
MemorySize(cacheMemory));
this.stateFactory =
- new RocksDBStateFactory(
+ new LocalKvStateFactory(
path.toString(),
- rocksdbOptions,
- coreOptions.crossPartitionUpsertIndexTtl());
+ stateOptions,
+ coreOptions.crossPartitionUpsertIndexTtl(),
+ compactionExecutor,
+ true);
RowType keyType = table.schema().logicalPrimaryKeysType();
this.keyIndex =
stateFactory.valueState(
@@ -167,7 +177,7 @@ public class GlobalIndexAssigner implements Serializable,
Closeable {
// create bootstrap sort buffer
this.bootstrap = true;
- this.bootstrapKeys = RocksDBState.createBulkLoadSorter(ioManager,
coreOptions);
+ this.bootstrapKeys = StateUtils.createBulkLoadSorter(ioManager,
coreOptions);
this.bootstrapRecords =
RowBuffer.getBuffer(
ioManager,
@@ -208,14 +218,14 @@ public class GlobalIndexAssigner implements Serializable,
Closeable {
bootstrap = false;
boolean isEmpty = true;
if (!bootstrapKeys.isEmpty()) {
- RocksDBBulkLoader bulkLoader = keyIndex.createBulkLoader();
+ ValueBulkLoader bulkLoader = keyIndex.createBulkLoader();
MutableObjectIterator<BinaryRow> keyIterator =
bootstrapKeys.sortedIterator();
BinaryRow row = new BinaryRow(2);
try {
while ((row = keyIterator.next(row)) != null) {
bulkLoader.write(row.getBinary(0), row.getBinary(1));
}
- } catch (RocksDBBulkLoader.WriteException e) {
+ } catch (BulkLoader.WriteException e) {
throw new RuntimeException(
"Exception in bulkLoad, the most suspicious reason is
that "
+ "your data contains duplicates, please check
your sink table. "
@@ -285,7 +295,7 @@ public class GlobalIndexAssigner implements Serializable,
Closeable {
// ================== End Public API ===================
- /** Sort bootstrap records and assign bucket without RocksDB. */
+ /** Sort bootstrap records and assign buckets without state lookups. */
private void bulkLoadBootstrapRecords() {
RowType rowType = table.rowType();
List<DataType> fields =
diff --git
a/paimon-core/src/main/java/org/apache/paimon/lookup/local/LocalKvListMergeOperator.java
b/paimon-core/src/main/java/org/apache/paimon/lookup/local/LocalKvListMergeOperator.java
index 7d9f33107f..f2b7f5fd8f 100644
---
a/paimon-core/src/main/java/org/apache/paimon/lookup/local/LocalKvListMergeOperator.java
+++
b/paimon-core/src/main/java/org/apache/paimon/lookup/local/LocalKvListMergeOperator.java
@@ -55,11 +55,6 @@ final class LocalKvListMergeOperator implements
LocalKvDb.MergeOperator {
return true;
}
- @Override
- public boolean canMergeTombstone(MemorySlice firstKey, MemorySlice
tombstoneKey) {
- return canMerge(firstKey, tombstoneKey);
- }
-
@Override
public byte[] merge(List<byte[]> values) throws IOException {
// Foreground flush and background compaction can invoke the operator
concurrently.
diff --git
a/paimon-core/src/main/java/org/apache/paimon/lookup/local/LocalKvListValueCodec.java
b/paimon-core/src/main/java/org/apache/paimon/lookup/local/LocalKvListValueCodec.java
index 3e238bb46b..3fec26f291 100644
---
a/paimon-core/src/main/java/org/apache/paimon/lookup/local/LocalKvListValueCodec.java
+++
b/paimon-core/src/main/java/org/apache/paimon/lookup/local/LocalKvListValueCodec.java
@@ -53,36 +53,51 @@ final class LocalKvListValueCodec {
}
byte[] merge(List<byte[]> storedValues, LocalKvValueCodec valueCodec)
throws IOException {
- long[] stats = new long[2];
+ long valueCount = 0;
+ long payloadLength = 0;
for (byte[] stored : storedValues) {
- inspectStoredValue(stored, valueCodec, stats);
+ int valueOffset = valueCodec.valueOffset(stored, 0, stored.length);
+ input.setBuffer(stored, valueOffset, stored.length - valueOffset);
+ int type = input.readUnsignedByte();
+ if (type == SINGLE_VALUE) {
+ valueCount++;
+ payloadLength += Integer.BYTES + (long) input.available();
+ continue;
+ }
+ if (type != PACKED_VALUES) {
+ throw new IOException("Corrupted local KV list value marker.");
+ }
+
+ int size = input.readInt();
+ if (size < 0 || size > input.available() / Integer.BYTES) {
+ throw new IOException("Corrupted local KV list size: " + size
+ '.');
+ }
+ int storedPayloadLength = input.available();
+ for (int i = 0; i < size; i++) {
+ int elementLength = input.readInt();
+ if (elementLength < 0 || elementLength > input.available()) {
+ throw new IOException(
+ "Corrupted local KV list element length: " +
elementLength + '.');
+ }
+ input.skipBytesToRead(elementLength);
+ }
+ if (input.available() != 0) {
+ throw new IOException(
+ "Corrupted local KV list with " + input.available() +
" trailing bytes.");
+ }
+ valueCount += size;
+ payloadLength += storedPayloadLength;
}
- if (stats[0] > Integer.MAX_VALUE || stats[1] > Integer.MAX_VALUE - 5) {
+ if (valueCount > Integer.MAX_VALUE || payloadLength >
Integer.MAX_VALUE - 5) {
throw new IOException("Merged local KV list value is too large.");
}
- byte[] packed = new byte[5 + (int) stats[1]];
+ byte[] packed = new byte[5 + (int) payloadLength];
packed[0] = PACKED_VALUES;
- writeInt(packed, 1, (int) stats[0]);
+ writeInt(packed, 1, (int) valueCount);
int outputOffset = 5;
for (byte[] stored : storedValues) {
- int valueOffset = valueCodec.valueOffset(stored, 0, stored.length);
- input.setBuffer(stored, valueOffset, stored.length - valueOffset);
- int type = input.readUnsignedByte();
- if (type == SINGLE_VALUE) {
- int valueLength = input.available();
- writeInt(packed, outputOffset, valueLength);
- outputOffset += Integer.BYTES;
- System.arraycopy(stored, input.getPosition(), packed,
outputOffset, valueLength);
- outputOffset += valueLength;
- } else if (type == PACKED_VALUES) {
- input.readInt();
- int payloadLength = input.available();
- System.arraycopy(stored, input.getPosition(), packed,
outputOffset, payloadLength);
- outputOffset += payloadLength;
- } else {
- throw new IOException("Corrupted local KV list value marker.");
- }
+ outputOffset = copyStoredValue(stored, valueCodec, packed,
outputOffset);
}
return valueCodec.encode(packed);
}
@@ -143,39 +158,27 @@ final class LocalKvListValueCodec {
return value;
}
- private void inspectStoredValue(byte[] stored, LocalKvValueCodec
valueCodec, long[] stats)
+ private int copyStoredValue(
+ byte[] stored, LocalKvValueCodec valueCodec, byte[] target, int
targetOffset)
throws IOException {
int valueOffset = valueCodec.valueOffset(stored, 0, stored.length);
input.setBuffer(stored, valueOffset, stored.length - valueOffset);
int type = input.readUnsignedByte();
if (type == SINGLE_VALUE) {
- stats[0]++;
- stats[1] += Integer.BYTES + input.available();
- return;
+ int valueLength = input.available();
+ writeInt(target, targetOffset, valueLength);
+ targetOffset += Integer.BYTES;
+ System.arraycopy(stored, input.getPosition(), target,
targetOffset, valueLength);
+ return targetOffset + valueLength;
}
if (type != PACKED_VALUES) {
throw new IOException("Corrupted local KV list value marker.");
}
- int size = input.readInt();
- if (size < 0 || size > input.available() / Integer.BYTES) {
- throw new IOException("Corrupted local KV list size: " + size +
'.');
- }
+ input.readInt();
int payloadLength = input.available();
- for (int i = 0; i < size; i++) {
- int elementLength = input.readInt();
- if (elementLength < 0 || elementLength > input.available()) {
- throw new IOException(
- "Corrupted local KV list element length: " +
elementLength + '.');
- }
- input.skipBytesToRead(elementLength);
- }
- if (input.available() != 0) {
- throw new IOException(
- "Corrupted local KV list with " + input.available() + "
trailing bytes.");
- }
- stats[0] += size;
- stats[1] += payloadLength;
+ System.arraycopy(stored, input.getPosition(), target, targetOffset,
payloadLength);
+ return targetOffset + payloadLength;
}
private static void writeInt(byte[] bytes, int offset, int value) {
diff --git
a/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java
b/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java
index 65d916439b..821521d255 100644
---
a/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java
@@ -93,7 +93,7 @@ public class GlobalIndexAssignerTest extends TableTestBase {
private void innerTestBucketAssign(boolean enableTtl) throws Exception {
GlobalIndexAssigner assigner = createAssigner(MergeEngine.DEDUPLICATE,
enableTtl);
List<Integer> output = new ArrayList<>();
- assigner.open(0, ioManager(), 2, 0, (row, bucket) ->
output.add(bucket));
+ assigner.open(0, null, ioManager(), 2, 0, (row, bucket) ->
output.add(bucket));
assigner.endBoostrap(false);
// assign
@@ -127,7 +127,8 @@ public class GlobalIndexAssignerTest extends TableTestBase {
public void testUpsert() throws Exception {
GlobalIndexAssigner assigner = createAssigner(MergeEngine.DEDUPLICATE);
List<Pair<InternalRow, Integer>> output = new ArrayList<>();
- assigner.open(0, ioManager(), 2, 0, (row, bucket) ->
output.add(Pair.of(row, bucket)));
+ assigner.open(
+ 0, null, ioManager(), 2, 0, (row, bucket) ->
output.add(Pair.of(row, bucket)));
assigner.endBoostrap(false);
// change partition
@@ -171,7 +172,8 @@ public class GlobalIndexAssignerTest extends TableTestBase {
: MergeEngine.AGGREGATE;
GlobalIndexAssigner assigner = createAssigner(mergeEngine);
List<Pair<InternalRow, Integer>> output = new ArrayList<>();
- assigner.open(0, ioManager(), 2, 0, (row, bucket) ->
output.add(Pair.of(row, bucket)));
+ assigner.open(
+ 0, null, ioManager(), 2, 0, (row, bucket) ->
output.add(Pair.of(row, bucket)));
assigner.endBoostrap(false);
// change partition
@@ -195,7 +197,8 @@ public class GlobalIndexAssignerTest extends TableTestBase {
public void testFirstRow() throws Exception {
GlobalIndexAssigner assigner = createAssigner(MergeEngine.FIRST_ROW);
List<Pair<InternalRow, Integer>> output = new ArrayList<>();
- assigner.open(0, ioManager(), 2, 0, (row, bucket) ->
output.add(Pair.of(row, bucket)));
+ assigner.open(
+ 0, null, ioManager(), 2, 0, (row, bucket) ->
output.add(Pair.of(row, bucket)));
assigner.endBoostrap(false);
// change partition
@@ -219,6 +222,7 @@ public class GlobalIndexAssignerTest extends TableTestBase {
List<List<Integer>> output = new ArrayList<>();
assigner.open(
0,
+ null,
ioManager(),
2,
0,
@@ -251,6 +255,7 @@ public class GlobalIndexAssignerTest extends TableTestBase {
List<List<Integer>> output = new ArrayList<>();
assigner.open(
0,
+ null,
ioManager(),
2,
0,
diff --git
a/paimon-core/src/test/java/org/apache/paimon/lookup/local/LocalKvStateFactoryTest.java
b/paimon-core/src/test/java/org/apache/paimon/lookup/local/LocalKvStateFactoryTest.java
index e89fe40142..1612d7cc48 100644
---
a/paimon-core/src/test/java/org/apache/paimon/lookup/local/LocalKvStateFactoryTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/lookup/local/LocalKvStateFactoryTest.java
@@ -211,6 +211,54 @@ class LocalKvStateFactoryTest {
}
}
+ @Test
+ void testListStateCollectsMemTableDeltasLazily() throws Exception {
+ try (LocalKvStateFactory factory = createFactory()) {
+ @SuppressWarnings("unchecked")
+ LocalKvListState<Integer, Integer> state =
+ (LocalKvListState<Integer, Integer>)
+ factory.listState(
+ "batched-list",
+ IntSerializer.INSTANCE,
+ IntSerializer.INSTANCE,
+ 10);
+ List<Integer> firstExpected = new ArrayList<>();
+ List<Integer> secondExpected = new ArrayList<>();
+ for (int value = 0; value < 100; value++) {
+ state.add(1, value);
+ state.add(2, 1_000 + value);
+ firstExpected.add(value);
+ secondExpected.add(1_000 + value);
+ }
+ for (int value = 0; value < 32; value++) {
+ state.add(3, value);
+ }
+ for (int value = 0; value < 33; value++) {
+ state.add(4, value);
+ }
+
+ assertThat(rawEntries(state, 1)).hasSize(1);
+ assertThat(rawEntries(state, 2)).hasSize(1);
+ assertThat(rawEntries(state, 3)).hasSize(1);
+ assertThat(rawEntries(state, 4)).hasSize(1);
+ assertThat(state.get(1)).containsExactlyElementsOf(firstExpected);
+ assertThat(state.get(2)).containsExactlyElementsOf(secondExpected);
+
assertThat(state.get(3)).containsExactlyElementsOf(firstExpected.subList(0,
32));
+
assertThat(state.get(4)).containsExactlyElementsOf(firstExpected.subList(0,
33));
+
+ state.db.flush();
+ state.cache.invalidateAll();
+ assertThat(rawEntries(state, 1)).hasSize(1);
+ assertThat(rawEntries(state, 2)).hasSize(1);
+ assertThat(rawEntries(state, 3)).hasSize(1);
+ assertThat(rawEntries(state, 4)).hasSize(1);
+ assertThat(state.get(1)).containsExactlyElementsOf(firstExpected);
+ assertThat(state.get(2)).containsExactlyElementsOf(secondExpected);
+
assertThat(state.get(3)).containsExactlyElementsOf(firstExpected.subList(0,
32));
+
assertThat(state.get(4)).containsExactlyElementsOf(firstExpected.subList(0,
33));
+ }
+ }
+
@Test
void testListStateMergesFragmentsDuringFlushAndCompaction() throws
Exception {
try (LocalKvStateFactory factory = createFactory()) {
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkConnectorOptions.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkConnectorOptions.java
index 1f0bf32c09..2eb488dd7d 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkConnectorOptions.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/FlinkConnectorOptions.java
@@ -224,7 +224,7 @@ public class FlinkConnectorOptions {
.memoryType()
.defaultValue(MemorySize.ofMebiBytes(256))
.withDescription(
- "Weight of managed memory for RocksDB in
cross-partition update, Flink will compute the memory size "
+ "Weight of managed memory for the local key-value
index in cross-partition update, Flink will compute the memory size "
+ "according to the weight, the actual
memory used depends on the running environment.");
public static final ConfigOption<Boolean> SOURCE_CHECKPOINT_ALIGN_ENABLED =
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/FullCacheLookupTable.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/FullCacheLookupTable.java
index 937d8357ca..08cae53350 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/FullCacheLookupTable.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/FullCacheLookupTable.java
@@ -24,11 +24,11 @@ import org.apache.paimon.data.BinaryRow;
import org.apache.paimon.data.GenericRow;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.disk.IOManager;
+import org.apache.paimon.lookup.BulkLoader;
import org.apache.paimon.lookup.StateFactory;
+import org.apache.paimon.lookup.StateUtils;
+import org.apache.paimon.lookup.local.LocalKvStateFactory;
import org.apache.paimon.lookup.memory.InMemoryStateFactory;
-import org.apache.paimon.lookup.rocksdb.RocksDBBulkLoader;
-import org.apache.paimon.lookup.rocksdb.RocksDBState;
-import org.apache.paimon.lookup.rocksdb.RocksDBStateFactory;
import org.apache.paimon.options.Options;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
@@ -185,7 +185,7 @@ public abstract class FullCacheLookupTable implements
LookupTable {
if (options.get(LOOKUP_CACHE_MODE) == MEMORY) {
return new InMemoryStateFactory();
} else {
- return new RocksDBStateFactory(diskDir, options, null);
+ return new LocalKvStateFactory(diskDir, options, null, null,
false);
}
}
@@ -224,7 +224,7 @@ public abstract class FullCacheLookupTable implements
LookupTable {
boolean useParallelBootstrapRead = !(blobAsDescriptor &&
hasBlobFileFields);
BinaryExternalSortBuffer bulkLoadSorter =
- RocksDBState.createBulkLoadSorter(
+ StateUtils.createBulkLoadSorter(
IOManager.create(context.tempPath.toString()),
context.table.coreOptions());
Predicate predicate = projectedPredicate();
try (RecordReaderIterator<InternalRow> batch =
@@ -246,7 +246,7 @@ public abstract class FullCacheLookupTable implements
LookupTable {
while ((row = keyIterator.next(row)) != null) {
bulkLoader.write(row.getBinary(0), row.getBinary(1));
}
- } catch (RocksDBBulkLoader.WriteException e) {
+ } catch (BulkLoader.WriteException e) {
throw new RuntimeException(
"Exception in bulkLoad, the most suspicious reason is that
"
+ "your data contains duplicates, please check
your lookup table. ",
@@ -434,7 +434,7 @@ public abstract class FullCacheLookupTable implements
LookupTable {
/** Bulk loader for the table. */
public interface TableBulkLoader {
- void write(byte[] key, byte[] value) throws
RocksDBBulkLoader.WriteException, IOException;
+ void write(byte[] key, byte[] value) throws BulkLoader.WriteException,
IOException;
void finish() throws IOException;
}
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/NoPrimaryKeyLookupTable.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/NoPrimaryKeyLookupTable.java
index 4fd1dc1ae0..302b07d5ae 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/NoPrimaryKeyLookupTable.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/NoPrimaryKeyLookupTable.java
@@ -20,9 +20,9 @@ package org.apache.paimon.flink.lookup;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.serializer.InternalSerializers;
+import org.apache.paimon.lookup.BulkLoader;
import org.apache.paimon.lookup.ListBulkLoader;
import org.apache.paimon.lookup.ListState;
-import org.apache.paimon.lookup.rocksdb.RocksDBBulkLoader;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.types.RowKind;
import org.apache.paimon.utils.KeyProjectedRow;
@@ -132,7 +132,7 @@ public class NoPrimaryKeyLookupTable extends
FullCacheLookupTable {
if (currentKey != null && values.size() > 0) {
try {
bulkLoader.write(currentKey, values);
- } catch (RocksDBBulkLoader.WriteException e) {
+ } catch (BulkLoader.WriteException e) {
throw new RuntimeException(e);
}
}
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/PrimaryKeyLookupTable.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/PrimaryKeyLookupTable.java
index 485a9cfaa6..172d92cc94 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/PrimaryKeyLookupTable.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/PrimaryKeyLookupTable.java
@@ -20,9 +20,9 @@ package org.apache.paimon.flink.lookup;
import org.apache.paimon.data.InternalRow;
import org.apache.paimon.data.serializer.InternalSerializers;
+import org.apache.paimon.lookup.BulkLoader;
import org.apache.paimon.lookup.ValueBulkLoader;
import org.apache.paimon.lookup.ValueState;
-import org.apache.paimon.lookup.rocksdb.RocksDBBulkLoader;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.table.FileStoreTable;
import org.apache.paimon.types.RowKind;
@@ -135,7 +135,7 @@ public class PrimaryKeyLookupTable extends
FullCacheLookupTable {
@Override
public void write(byte[] key, byte[] value)
- throws RocksDBBulkLoader.WriteException, IOException {
+ throws BulkLoader.WriteException, IOException {
bulkLoader.write(key, value);
bulkLoadWritePlus(key, value);
}
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/index/GlobalDynamicBucketSink.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/index/GlobalDynamicBucketSink.java
index 52efda2517..aaf05e78e0 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/index/GlobalDynamicBucketSink.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/index/GlobalDynamicBucketSink.java
@@ -126,7 +126,7 @@ public class GlobalDynamicBucketSink extends
FlinkWriteSink<Tuple2<InternalRow,
GlobalIndexAssignerOperator.forRowData(table))
.setParallelism(partitionByKeyHash.getParallelism());
- // declare managed memory for RocksDB
+ // declare managed memory for the local key-value index
declareManagedMemory(
bucketAssigned,
options.toConfiguration().get(SINK_CROSS_PARTITION_MANAGED_MEMORY));
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/index/GlobalIndexAssignerOperator.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/index/GlobalIndexAssignerOperator.java
index 99cce07fdc..017d0f71a4 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/index/GlobalIndexAssignerOperator.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/index/GlobalIndexAssignerOperator.java
@@ -59,6 +59,7 @@ public class GlobalIndexAssignerOperator
ioManager =
IOManager.create(flinkIoManager.getSpillingDirectoriesPaths());
assigner.open(
computeManagedMemory(this),
+
getContainingTask().getEnvironment().getAsyncOperationsThreadPool(),
ioManager,
RuntimeContextUtils.getNumberOfParallelSubtasks(getRuntimeContext()),
RuntimeContextUtils.getIndexOfThisSubtask(getRuntimeContext()),
diff --git
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/CrossPartitionTableITCase.java
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/CrossPartitionTableITCase.java
index ddad7c5e62..8b217dcf17 100644
---
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/CrossPartitionTableITCase.java
+++
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/CrossPartitionTableITCase.java
@@ -157,7 +157,6 @@ public class CrossPartitionTableITCase extends
CatalogITCaseBase {
sql(
"create table large_t (pt int, k int, v int, primary key (k)
not enforced) partitioned by (pt) with ("
+ "'bucket'='-1', "
- + "'rocksdb.compaction.level.target-file-size-base'='2
kb', "
+ "'dynamic-bucket.target-row-num'='10000')");
sql(
"create temporary table src (pt int, k int, v int) with ("
diff --git
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupTableTest.java
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupTableTest.java
index b077e9dcdb..f3b291b976 100644
---
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupTableTest.java
+++
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/LookupTableTest.java
@@ -153,7 +153,7 @@ public class LookupTableTest extends TableTestBase {
TableBulkLoader bulkLoader = table.createBulkLoader();
bulkLoader.write(new byte[] {1}, new byte[] {1});
assertThatThrownBy(() -> bulkLoader.write(new byte[] {1}, new
byte[] {2}))
- .hasMessageContaining("Keys must be added in strict
ascending order");
+ .hasMessageContaining("strictly increasing");
}
// test bulk load 100_000 records
diff --git
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/BucketProcessor.scala
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/BucketProcessor.scala
index e6cb3b55ab..4a147deee9 100644
---
a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/BucketProcessor.scala
+++
b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/commands/BucketProcessor.scala
@@ -210,6 +210,7 @@ class GlobalIndexAssignerIterator(
val _assigner = new GlobalIndexAssigner(fileStoreTable)
_assigner.open(
0,
+ null,
ioManager,
numAssigners,
TaskContext.getPartitionId(),