From 47854125397be01f8926686a60c7701b2c533c78 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=99=BD=E9=B5=BA?= Date: Fri, 28 Aug 2026 10:53:25 +0800 Subject: [PATCH 1/6] [kv] Bump RocksDB version to 11.8.1 Replace FRocksDB with fluss-rocksdbjni 11.8.1-fluss-2 and adapt the JNI package and API changes. Handle empty checkpoint files generated by newer RocksDB releases. Closes #4132. --- .../scanner/batch/SnapshotFilesReader.java | 12 +++---- .../src/main/resources/META-INF/NOTICE | 4 +-- fluss-common/pom.xml | 4 +-- .../apache/fluss/rocksdb/RocksDBHandle.java | 10 +++--- .../fluss/rocksdb/RocksDBOperationUtils.java | 12 +++---- .../fluss/rocksdb/RocksIteratorWrapper.java | 13 ++++++-- .../org/apache/fluss/server/kv/KvManager.java | 12 +++---- .../org/apache/fluss/server/kv/KvTablet.java | 16 +++++----- .../kv/RowTtlCompactionFilterFactory.java | 18 +++++------ .../fluss/server/kv/rocksdb/RocksDBKv.java | 22 ++++++------- .../server/kv/rocksdb/RocksDBKvBuilder.java | 14 ++++---- .../kv/rocksdb/RocksDBResourceContainer.java | 32 +++++++++---------- .../server/kv/rocksdb/RocksDBStatistics.java | 20 ++++++------ .../kv/rocksdb/RocksDBWriteBatchWrapper.java | 10 +++--- .../fluss/server/kv/scan/ScannerContext.java | 8 ++--- .../kv/snapshot/KvSnapshotDataUploader.java | 2 ++ .../kv/snapshot/RocksIncrementalSnapshot.java | 4 +-- .../src/main/resources/META-INF/NOTICE | 3 +- .../apache/fluss/server/kv/KvTabletTest.java | 6 ++-- .../server/kv/RowTtlCompactionFilterTest.java | 12 +++---- .../server/kv/rocksdb/RocksDBExtension.java | 2 +- .../kv/rocksdb/RocksDBKvBuilderTest.java | 6 ++-- .../server/kv/rocksdb/RocksDBKvTest.java | 8 ++--- .../server/kv/rocksdb/RocksDBKvTestUtils.java | 2 +- .../rocksdb/RocksDBOperationsUtilsTest.java | 13 +++++--- .../rocksdb/RocksDBResourceContainerTest.java | 30 ++++++++--------- .../snapshot/KvTabletSnapshotTargetTest.java | 2 +- .../RocksIncrementalSnapshotTest.java | 2 +- .../fluss/server/testutils/KvTestUtils.java | 4 +-- pom.xml | 8 ++--- 30 files changed, 162 insertions(+), 149 deletions(-) diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SnapshotFilesReader.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SnapshotFilesReader.java index 5a6cbac5a1f..65f823bf87f 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SnapshotFilesReader.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SnapshotFilesReader.java @@ -34,12 +34,12 @@ import org.apache.fluss.utils.IOUtils; import org.apache.fluss.utils.SchemaUtil; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.DBOptions; -import org.rocksdb.ReadOptions; -import org.rocksdb.RocksDB; -import org.rocksdb.RocksIterator; -import org.rocksdb.Snapshot; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.DBOptions; +import org.fluss.rocksdb.ReadOptions; +import org.fluss.rocksdb.RocksDB; +import org.fluss.rocksdb.RocksIterator; +import org.fluss.rocksdb.Snapshot; import javax.annotation.Nullable; import javax.annotation.concurrent.NotThreadSafe; diff --git a/fluss-client/src/main/resources/META-INF/NOTICE b/fluss-client/src/main/resources/META-INF/NOTICE index d8059b65de5..6f1fe5166ff 100644 --- a/fluss-client/src/main/resources/META-INF/NOTICE +++ b/fluss-client/src/main/resources/META-INF/NOTICE @@ -7,7 +7,7 @@ The Apache Software Foundation (http://www.apache.org/). This project bundles the following dependencies under the Apache Software License 2.0 (http://www.apache.org/licenses/LICENSE-2.0.txt) - com.google.code.findbugs:jsr305:1.3.9 -- com.ververica:frocksdbjni:6.20.3-ververica-2.0 +- io.github.fluss-contrib:fluss-rocksdbjni:11.8.1-fluss-2 - org.apache.commons:commons-lang3:3.18.0 - org.apache.commons:commons-math3:3.6.1 - at.yawk.lz4:lz4-java:1.10.2 @@ -20,4 +20,4 @@ See bundled license files for details. This project bundles the following dependencies under BSD License (https://opensource.org/licenses/bsd-license.php). See bundled license files for details. -- com.github.luben:zstd-jni:1.5.7-6 \ No newline at end of file +- com.github.luben:zstd-jni:1.5.7-6 diff --git a/fluss-common/pom.xml b/fluss-common/pom.xml index e711b12b203..8f54c94824e 100644 --- a/fluss-common/pom.xml +++ b/fluss-common/pom.xml @@ -98,8 +98,8 @@ the rocksdb should be provided as a kv plugin to used by client & server. --> - com.ververica - frocksdbjni + io.github.fluss-contrib + fluss-rocksdbjni diff --git a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBHandle.java b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBHandle.java index 823f2723494..6b81747923d 100644 --- a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBHandle.java +++ b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBHandle.java @@ -19,11 +19,11 @@ import org.apache.fluss.utils.IOUtils; -import org.rocksdb.ColumnFamilyDescriptor; -import org.rocksdb.ColumnFamilyHandle; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.DBOptions; -import org.rocksdb.RocksDB; +import org.fluss.rocksdb.ColumnFamilyDescriptor; +import org.fluss.rocksdb.ColumnFamilyHandle; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.DBOptions; +import org.fluss.rocksdb.RocksDB; import java.io.File; import java.io.IOException; diff --git a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBOperationUtils.java b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBOperationUtils.java index 8a9edff75a5..4060cfe37c5 100644 --- a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBOperationUtils.java +++ b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBOperationUtils.java @@ -20,12 +20,12 @@ import org.apache.fluss.utils.IOUtils; import org.apache.fluss.utils.OperatingSystem; -import org.rocksdb.ColumnFamilyDescriptor; -import org.rocksdb.ColumnFamilyHandle; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.DBOptions; -import org.rocksdb.RocksDB; -import org.rocksdb.RocksDBException; +import org.fluss.rocksdb.ColumnFamilyDescriptor; +import org.fluss.rocksdb.ColumnFamilyHandle; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.DBOptions; +import org.fluss.rocksdb.RocksDB; +import org.fluss.rocksdb.RocksDBException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksIteratorWrapper.java b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksIteratorWrapper.java index 1d1591dfeb3..849964ffb31 100644 --- a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksIteratorWrapper.java +++ b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksIteratorWrapper.java @@ -17,9 +17,10 @@ package org.apache.fluss.rocksdb; -import org.rocksdb.RocksDBException; -import org.rocksdb.RocksIterator; -import org.rocksdb.RocksIteratorInterface; +import org.fluss.rocksdb.RocksDBException; +import org.fluss.rocksdb.RocksIterator; +import org.fluss.rocksdb.RocksIteratorInterface; +import org.fluss.rocksdb.Snapshot; import javax.annotation.Nonnull; @@ -114,6 +115,12 @@ public void refresh() throws RocksDBException { status(); } + @Override + public void refresh(Snapshot snapshot) throws RocksDBException { + iterator.refresh(snapshot); + status(); + } + public byte[] key() { return iterator.key(); } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java index 29eedc5eb6c..f3a17bf6a0d 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java @@ -54,12 +54,12 @@ import org.apache.fluss.utils.function.SupplierWithException; import org.apache.fluss.utils.types.Tuple2; -import org.rocksdb.Cache; -import org.rocksdb.LRUCache; -import org.rocksdb.RateLimiter; -import org.rocksdb.RateLimiterMode; -import org.rocksdb.RocksDB; -import org.rocksdb.WriteBufferManager; +import org.fluss.rocksdb.Cache; +import org.fluss.rocksdb.LRUCache; +import org.fluss.rocksdb.RateLimiter; +import org.fluss.rocksdb.RateLimiterMode; +import org.fluss.rocksdb.RocksDB; +import org.fluss.rocksdb.WriteBufferManager; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java index e6bea68f7ab..d39df7a7787 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java @@ -70,14 +70,14 @@ import org.apache.fluss.utils.clock.Clock; import org.apache.fluss.utils.clock.SystemClock; -import org.rocksdb.AbstractCompactionFilter; -import org.rocksdb.AbstractCompactionFilterFactory; -import org.rocksdb.Cache; -import org.rocksdb.RateLimiter; -import org.rocksdb.ReadOptions; -import org.rocksdb.RocksIterator; -import org.rocksdb.Snapshot; -import org.rocksdb.WriteBufferManager; +import org.fluss.rocksdb.AbstractCompactionFilter; +import org.fluss.rocksdb.AbstractCompactionFilterFactory; +import org.fluss.rocksdb.Cache; +import org.fluss.rocksdb.RateLimiter; +import org.fluss.rocksdb.ReadOptions; +import org.fluss.rocksdb.RocksIterator; +import org.fluss.rocksdb.Snapshot; +import org.fluss.rocksdb.WriteBufferManager; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/RowTtlCompactionFilterFactory.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/RowTtlCompactionFilterFactory.java index e77e5b9b468..2c1698cf33a 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/RowTtlCompactionFilterFactory.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/RowTtlCompactionFilterFactory.java @@ -21,8 +21,8 @@ import org.apache.fluss.server.utils.RowTtlUtils; import org.apache.fluss.utils.clock.Clock; -import org.rocksdb.FlinkCompactionFilter; -import org.rocksdb.RocksDB; +import org.fluss.rocksdb.FlussTtlCompactionFilter; +import org.fluss.rocksdb.RocksDB; import java.time.Duration; import java.util.function.LongSupplier; @@ -38,7 +38,7 @@ public final class RowTtlCompactionFilterFactory { private RowTtlCompactionFilterFactory() {} /** Creates a configured native compaction filter factory for row TTL cleanup. */ - public static FlinkCompactionFilter.FlinkCompactionFilterFactory create( + public static FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory create( KvValueLayout kvValueLayout, Duration ttl, Clock clock) { long ttlMillis = RowTtlUtils.validateAndConvertTtlDurationToMillis(ttl); checkNotNull(clock, "clock must not be null."); @@ -46,7 +46,7 @@ public static FlinkCompactionFilter.FlinkCompactionFilterFactory create( } /** Removes values using the default interval for refreshing the supplied current value. */ - static FlinkCompactionFilter.FlinkCompactionFilterFactory create( + static FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory create( KvValueLayout kvValueLayout, long expirationDistance, LongSupplier currentValueSupplier) { @@ -58,7 +58,7 @@ static FlinkCompactionFilter.FlinkCompactionFilterFactory create( } /** Removes a value when {@code valueTag + expirationDistance <= currentValue}. */ - static FlinkCompactionFilter.FlinkCompactionFilterFactory create( + static FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory create( KvValueLayout kvValueLayout, long expirationDistance, long queryCurrentValueAfterNumEntries, @@ -72,12 +72,12 @@ static FlinkCompactionFilter.FlinkCompactionFilterFactory create( "queryCurrentValueAfterNumEntries must be greater than zero."); RocksDB.loadLibrary(); - FlinkCompactionFilter.FlinkCompactionFilterFactory factory = - new FlinkCompactionFilter.FlinkCompactionFilterFactory( + FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory factory = + new FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory( currentValueSupplier::getAsLong); factory.configure( - FlinkCompactionFilter.Config.createNotList( - FlinkCompactionFilter.StateType.Value, + FlussTtlCompactionFilter.Config.createNotList( + FlussTtlCompactionFilter.StateType.Value, kvValueLayout.valueTagOffset(), expirationDistance, queryCurrentValueAfterNumEntries)); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java index ce8ef903c5c..fe3537da139 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java @@ -27,16 +27,16 @@ import org.apache.fluss.utils.BytesUtils; import org.apache.fluss.utils.IOUtils; -import org.rocksdb.Cache; -import org.rocksdb.ColumnFamilyHandle; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.MutableDBOptions; -import org.rocksdb.ReadOptions; -import org.rocksdb.RocksDB; -import org.rocksdb.RocksDBException; -import org.rocksdb.RocksIterator; -import org.rocksdb.Statistics; -import org.rocksdb.WriteOptions; +import org.fluss.rocksdb.Cache; +import org.fluss.rocksdb.ColumnFamilyHandle; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.MutableDBOptions; +import org.fluss.rocksdb.ReadOptions; +import org.fluss.rocksdb.RocksDB; +import org.fluss.rocksdb.RocksDBException; +import org.fluss.rocksdb.RocksIterator; +import org.fluss.rocksdb.Statistics; +import org.fluss.rocksdb.WriteOptions; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -49,7 +49,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; -/** A wrapper for the operation of {@link org.rocksdb.RocksDB}. */ +/** A wrapper for the operation of {@link org.fluss.rocksdb.RocksDB}. */ public class RocksDBKv implements AutoCloseable { private static final Logger LOG = LoggerFactory.getLogger(RocksDBKv.class); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java index b44e892829c..aadab2a54b5 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java @@ -24,12 +24,12 @@ import org.apache.fluss.utils.FileUtils; import org.apache.fluss.utils.IOUtils; -import org.rocksdb.AbstractCompactionFilter; -import org.rocksdb.AbstractCompactionFilterFactory; -import org.rocksdb.ColumnFamilyHandle; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.NativeLibraryLoader; -import org.rocksdb.RocksDB; +import org.fluss.rocksdb.AbstractCompactionFilter; +import org.fluss.rocksdb.AbstractCompactionFilterFactory; +import org.fluss.rocksdb.ColumnFamilyHandle; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.NativeLibraryLoader; +import org.fluss.rocksdb.RocksDB; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -265,7 +265,7 @@ static void ensureRocksDBIsLoaded( @VisibleForTesting static void resetRocksDBLoadedFlag() throws Exception { final Field initField = - org.rocksdb.NativeLibraryLoader.class.getDeclaredField("initialized"); + org.fluss.rocksdb.NativeLibraryLoader.class.getDeclaredField("initialized"); initField.setAccessible(true); initField.setBoolean(null, false); } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java index 627175755a5..c410d1bb337 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java @@ -26,22 +26,22 @@ import org.apache.fluss.utils.FileUtils; import org.apache.fluss.utils.IOUtils; -import org.rocksdb.BlockBasedTableConfig; -import org.rocksdb.BloomFilter; -import org.rocksdb.Cache; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.CompactionStyle; -import org.rocksdb.CompressionType; -import org.rocksdb.DBOptions; -import org.rocksdb.InfoLogLevel; -import org.rocksdb.LRUCache; -import org.rocksdb.PlainTableConfig; -import org.rocksdb.RateLimiter; -import org.rocksdb.ReadOptions; -import org.rocksdb.Statistics; -import org.rocksdb.TableFormatConfig; -import org.rocksdb.WriteBufferManager; -import org.rocksdb.WriteOptions; +import org.fluss.rocksdb.BlockBasedTableConfig; +import org.fluss.rocksdb.BloomFilter; +import org.fluss.rocksdb.Cache; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.CompactionStyle; +import org.fluss.rocksdb.CompressionType; +import org.fluss.rocksdb.DBOptions; +import org.fluss.rocksdb.InfoLogLevel; +import org.fluss.rocksdb.LRUCache; +import org.fluss.rocksdb.PlainTableConfig; +import org.fluss.rocksdb.RateLimiter; +import org.fluss.rocksdb.ReadOptions; +import org.fluss.rocksdb.Statistics; +import org.fluss.rocksdb.TableFormatConfig; +import org.fluss.rocksdb.WriteBufferManager; +import org.fluss.rocksdb.WriteOptions; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBStatistics.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBStatistics.java index 58fe5eaa4f9..dc1407a172d 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBStatistics.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBStatistics.java @@ -19,16 +19,16 @@ import org.apache.fluss.server.utils.ResourceGuard; -import org.rocksdb.Cache; -import org.rocksdb.ColumnFamilyHandle; -import org.rocksdb.HistogramData; -import org.rocksdb.HistogramType; -import org.rocksdb.MemoryUsageType; -import org.rocksdb.MemoryUtil; -import org.rocksdb.RocksDB; -import org.rocksdb.RocksDBException; -import org.rocksdb.Statistics; -import org.rocksdb.TickerType; +import org.fluss.rocksdb.Cache; +import org.fluss.rocksdb.ColumnFamilyHandle; +import org.fluss.rocksdb.HistogramData; +import org.fluss.rocksdb.HistogramType; +import org.fluss.rocksdb.MemoryUsageType; +import org.fluss.rocksdb.MemoryUtil; +import org.fluss.rocksdb.RocksDB; +import org.fluss.rocksdb.RocksDBException; +import org.fluss.rocksdb.Statistics; +import org.fluss.rocksdb.TickerType; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBWriteBatchWrapper.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBWriteBatchWrapper.java index 7aaeab059c9..cf34a2c1831 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBWriteBatchWrapper.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBWriteBatchWrapper.java @@ -24,11 +24,11 @@ import org.apache.fluss.server.kv.KvBatchWriter; import org.apache.fluss.utils.IOUtils; -import org.rocksdb.RocksDB; -import org.rocksdb.RocksDBException; -import org.rocksdb.Status; -import org.rocksdb.WriteBatch; -import org.rocksdb.WriteOptions; +import org.fluss.rocksdb.RocksDB; +import org.fluss.rocksdb.RocksDBException; +import org.fluss.rocksdb.Status; +import org.fluss.rocksdb.WriteBatch; +import org.fluss.rocksdb.WriteOptions; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/scan/ScannerContext.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/scan/ScannerContext.java index aeb49886c9c..15201cf28f9 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/scan/ScannerContext.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/scan/ScannerContext.java @@ -24,10 +24,10 @@ import org.apache.fluss.server.utils.ResourceGuard; import org.apache.fluss.utils.IOUtils; -import org.rocksdb.ReadOptions; -import org.rocksdb.RocksDBException; -import org.rocksdb.RocksIterator; -import org.rocksdb.Snapshot; +import org.fluss.rocksdb.ReadOptions; +import org.fluss.rocksdb.RocksDBException; +import org.fluss.rocksdb.RocksIterator; +import org.fluss.rocksdb.Snapshot; import javax.annotation.concurrent.NotThreadSafe; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/KvSnapshotDataUploader.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/KvSnapshotDataUploader.java index e0f958ec6ac..63907adab13 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/KvSnapshotDataUploader.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/KvSnapshotDataUploader.java @@ -150,6 +150,8 @@ private KvFileHandleAndLocalPath uploadLocalFileToSnapshotLocation( uploadedBytes += numBytes; } + outputStream.flushToFile(); + final KvFileHandle result; if (closeableRegistry.unregisterCloseable(outputStream)) { result = outputStream.closeAndGetHandle(); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshot.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshot.java index 9cde2d5d8de..56e28fdf30f 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshot.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshot.java @@ -23,8 +23,8 @@ import org.apache.fluss.utils.ExceptionUtils; import org.apache.fluss.utils.FileUtils; -import org.rocksdb.Checkpoint; -import org.rocksdb.RocksDB; +import org.fluss.rocksdb.Checkpoint; +import org.fluss.rocksdb.RocksDB; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/resources/META-INF/NOTICE b/fluss-server/src/main/resources/META-INF/NOTICE index 2604b80cb38..a6c20fe7b08 100644 --- a/fluss-server/src/main/resources/META-INF/NOTICE +++ b/fluss-server/src/main/resources/META-INF/NOTICE @@ -9,7 +9,7 @@ This project bundles the following dependencies under the Apache Software Licens - com.github.ben-manes.caffeine:caffeine:2.9.3 - com.google.code.findbugs:jsr305:1.3.9 - com.google.errorprone:error_prone_annotations:2.10.0 -- com.ververica:frocksdbjni:6.20.3-ververica-2.0 +- io.github.fluss-contrib:fluss-rocksdbjni:11.8.1-fluss-2 - commons-cli:commons-cli:1.5.0 - org.apache.commons:commons-lang3:3.18.0 - org.apache.commons:commons-math3:3.6.1 @@ -28,4 +28,3 @@ See bundled license files for details. - com.github.luben:zstd-jni:1.5.7-6 - diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvTabletTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvTabletTest.java index 568d9f9fa57..8a9d2752a49 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvTabletTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvTabletTest.java @@ -87,14 +87,14 @@ import org.apache.fluss.utils.clock.SystemClock; import org.apache.fluss.utils.concurrent.FlussScheduler; +import org.fluss.rocksdb.FlushOptions; +import org.fluss.rocksdb.RocksDBException; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; -import org.rocksdb.FlushOptions; -import org.rocksdb.RocksDBException; import javax.annotation.Nullable; @@ -2285,7 +2285,7 @@ void testRocksDBMetrics() throws Exception { assertThat(statistics).as("RocksDB statistics should be available").isNotNull(); // Verify statistics is properly initialized - org.rocksdb.Statistics stats = kvTablet.getRocksDBKv().getStatistics(); + org.fluss.rocksdb.Statistics stats = kvTablet.getRocksDBKv().getStatistics(); assertThat(stats).as("RocksDB Statistics should be enabled").isNotNull(); // All metrics should start at 0 for a fresh database diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/RowTtlCompactionFilterTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/RowTtlCompactionFilterTest.java index 644922c6d5a..c5e83bd0209 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/RowTtlCompactionFilterTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/RowTtlCompactionFilterTest.java @@ -25,12 +25,12 @@ import org.apache.fluss.row.encode.ValueEncoder; import org.apache.fluss.utils.clock.ManualClock; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.DBOptions; +import org.fluss.rocksdb.FlushOptions; +import org.fluss.rocksdb.FlussTtlCompactionFilter; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.DBOptions; -import org.rocksdb.FlinkCompactionFilter; -import org.rocksdb.FlushOptions; import java.nio.charset.StandardCharsets; import java.nio.file.Path; @@ -48,13 +48,13 @@ class RowTtlCompactionFilterTest { @TempDir private Path tempDir; @Test - void testFlinkCompactionFilterReadsTimestampFromTaggedValue() throws Exception { + void testFlussTtlCompactionFilterReadsTimestampFromTaggedValue() throws Exception { byte[] expiredKey = "expired-key".getBytes(StandardCharsets.UTF_8); byte[] freshKey = "fresh-key".getBytes(StandardCharsets.UTF_8); BinaryRow row = compactedRow(DATA1_ROW_TYPE, new Object[] {1, "a"}); long now = 123456789L; - try (FlinkCompactionFilter.FlinkCompactionFilterFactory filterFactory = + try (FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory filterFactory = RowTtlCompactionFilterFactory.create( KvValueLayout.TAGGED, Duration.ofHours(1L), new ManualClock(now)); DBOptions dbOptions = new DBOptions().setCreateIfMissing(true); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBExtension.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBExtension.java index 646aacf7791..cf3fcf78d28 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBExtension.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBExtension.java @@ -20,10 +20,10 @@ import org.apache.fluss.config.Configuration; import org.apache.fluss.utils.IOUtils; +import org.fluss.rocksdb.RocksDB; import org.junit.jupiter.api.extension.AfterEachCallback; import org.junit.jupiter.api.extension.BeforeEachCallback; import org.junit.jupiter.api.extension.ExtensionContext; -import org.rocksdb.RocksDB; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java index 7261fe4362f..23be6d7c8e0 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java @@ -24,9 +24,9 @@ import org.apache.fluss.server.kv.RowTtlCompactionFilterFactory; import org.apache.fluss.utils.clock.ManualClock; +import org.fluss.rocksdb.FlussTtlCompactionFilter; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; -import org.rocksdb.FlinkCompactionFilter; import java.io.File; import java.io.IOException; @@ -66,7 +66,7 @@ void testTempLibFolderDeletedOnFail(@TempDir Path tempDir) { @Test void testCompactionFilterFactoryClosedWithKv(@TempDir Path tempDir) throws Exception { - FlinkCompactionFilter.FlinkCompactionFilterFactory filterFactory = + FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory filterFactory = RowTtlCompactionFilterFactory.create( KvValueLayout.TAGGED, Duration.ofHours(1L), new ManualClock(0L)); RocksDBResourceContainer rocksDBResourceContainer = @@ -90,7 +90,7 @@ void testCompactionFilterFactoryClosedWithKv(@TempDir Path tempDir) throws Excep @Test void testCompactionFilterFactoryClosedWhenBuildFails() { - FlinkCompactionFilter.FlinkCompactionFilterFactory filterFactory = + FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory filterFactory = RowTtlCompactionFilterFactory.create( KvValueLayout.TAGGED, Duration.ofHours(1L), new ManualClock(0L)); RocksDBResourceContainer rocksDBResourceContainer = new RocksDBResourceContainer(); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java index 7add3604701..ac05209afa3 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java @@ -23,11 +23,11 @@ import org.apache.fluss.metrics.util.TestHistogram; import org.apache.fluss.server.kv.KvCloseMode; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.FlushOptions; +import org.fluss.rocksdb.RocksDBException; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.FlushOptions; -import org.rocksdb.RocksDBException; import java.io.File; import java.nio.file.Path; @@ -130,7 +130,7 @@ void testNoSlowdownWriteBatchFailsFastOnRocksDbDelay(@TempDir Path tempDir) thro } /** - * Verifies that the L0 property is actually readable on frocksdbjni 6.20.3-ververica-2.0. + * Verifies that the L0 property is actually readable on fluss-rocksdbjni 11.8.1-fluss-2. * *

This is a regression test: {@code getLongProperty("rocksdb.num-files-at-level0")} throws * {@code RocksDBException: NotFound} because the property is parametric (string-type), not an diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTestUtils.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTestUtils.java index 74e5f466188..75646ec414f 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTestUtils.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTestUtils.java @@ -18,7 +18,7 @@ package org.apache.fluss.server.kv.rocksdb; -import org.rocksdb.RocksDBException; +import org.fluss.rocksdb.RocksDBException; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.spy; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBOperationsUtilsTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBOperationsUtilsTest.java index d2f9d8b8268..adbb8a515c8 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBOperationsUtilsTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBOperationsUtilsTest.java @@ -19,11 +19,13 @@ import org.apache.fluss.rocksdb.RocksDBOperationUtils; +import org.fluss.rocksdb.ColumnFamilyDescriptor; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.DBOptions; +import org.fluss.rocksdb.RocksDB; +import org.fluss.rocksdb.RocksDBException; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; -import org.rocksdb.DBOptions; -import org.rocksdb.RocksDB; -import org.rocksdb.RocksDBException; import java.io.File; import java.io.IOException; @@ -48,7 +50,10 @@ void testOpenDBFail(@TempDir Path temporaryFolder) throws Exception { RocksDB rocks = RocksDBOperationUtils.openDB( rocksDir.getAbsolutePath(), - Collections.emptyList(), + Collections.singletonList( + new ColumnFamilyDescriptor( + RocksDB.DEFAULT_COLUMN_FAMILY, + new ColumnFamilyOptions())), Collections.emptyList(), dbOptions, false); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java index 013ac3d626e..2dd117d736d 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java @@ -21,23 +21,23 @@ import org.apache.fluss.config.Configuration; import org.apache.fluss.server.kv.KvManager; +import org.fluss.rocksdb.BlockBasedTableConfig; +import org.fluss.rocksdb.BloomFilter; +import org.fluss.rocksdb.Cache; +import org.fluss.rocksdb.ColumnFamilyOptions; +import org.fluss.rocksdb.CompactionStyle; +import org.fluss.rocksdb.CompressionType; +import org.fluss.rocksdb.DBOptions; +import org.fluss.rocksdb.FlushOptions; +import org.fluss.rocksdb.InfoLogLevel; +import org.fluss.rocksdb.LRUCache; +import org.fluss.rocksdb.RateLimiter; +import org.fluss.rocksdb.ReadOptions; +import org.fluss.rocksdb.WriteBufferManager; +import org.fluss.rocksdb.WriteOptions; +import org.fluss.rocksdb.util.SizeUnit; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; -import org.rocksdb.BlockBasedTableConfig; -import org.rocksdb.BloomFilter; -import org.rocksdb.Cache; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.CompactionStyle; -import org.rocksdb.CompressionType; -import org.rocksdb.DBOptions; -import org.rocksdb.FlushOptions; -import org.rocksdb.InfoLogLevel; -import org.rocksdb.LRUCache; -import org.rocksdb.RateLimiter; -import org.rocksdb.ReadOptions; -import org.rocksdb.WriteBufferManager; -import org.rocksdb.WriteOptions; -import org.rocksdb.util.SizeUnit; import java.io.File; import java.nio.file.Path; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/KvTabletSnapshotTargetTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/KvTabletSnapshotTargetTest.java index 00ab5dfd83d..3f1124488a8 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/KvTabletSnapshotTargetTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/KvTabletSnapshotTargetTest.java @@ -42,13 +42,13 @@ import org.apache.fluss.utils.FlussPaths; import org.apache.fluss.utils.concurrent.Executors; +import org.fluss.rocksdb.RocksDB; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; import org.junit.jupiter.api.io.TempDir; -import org.rocksdb.RocksDB; import javax.annotation.Nonnull; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java index 7446d950975..0d9b8b04a51 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java @@ -28,12 +28,12 @@ import org.apache.fluss.utils.CloseableRegistry; import org.apache.fluss.utils.FlussPaths; +import org.fluss.rocksdb.RocksDB; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; import org.junit.jupiter.api.io.TempDir; -import org.rocksdb.RocksDB; import java.nio.file.Path; import java.util.Collection; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/testutils/KvTestUtils.java b/fluss-server/src/test/java/org/apache/fluss/server/testutils/KvTestUtils.java index 0e9de09dd76..e735e681d65 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/testutils/KvTestUtils.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/testutils/KvTestUtils.java @@ -39,8 +39,8 @@ import org.apache.fluss.utils.FileUtils; import org.apache.fluss.utils.types.Tuple2; -import org.rocksdb.RocksDB; -import org.rocksdb.RocksIterator; +import org.fluss.rocksdb.RocksDB; +import org.fluss.rocksdb.RocksIterator; import javax.annotation.Nullable; diff --git a/pom.xml b/pom.xml index 940ca3456ec..d6367b5a0fa 100644 --- a/pom.xml +++ b/pom.xml @@ -107,7 +107,7 @@ 3.1.0 3.4.3 - 6.20.3-ververica-2.0 + 11.8.1-fluss-2 1.7.36 2.25.4 2.3.1 @@ -367,9 +367,9 @@ - com.ververica - frocksdbjni - ${frocksdb.version} + io.github.fluss-contrib + fluss-rocksdbjni + ${fluss-rocksdb.version} From 8a4ead702bf8aabfc8e934cc00570383fc248450 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=99=BD=E9=B5=BA?= Date: Fri, 28 Aug 2026 15:15:45 +0800 Subject: [PATCH 2/6] [kv] Verify snapshot compatibility with FRocksDB --- fluss-server/pom.xml | 9 ++- .../kv/rocksdb/RocksDBResourceContainer.java | 8 +++ .../kv/snapshot/FrocksDBSnapshotReader.java | 53 +++++++++++++++ .../RocksIncrementalSnapshotTest.java | 67 +++++++++++++++++++ pom.xml | 1 + 5 files changed, 137 insertions(+), 1 deletion(-) create mode 100644 fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/FrocksDBSnapshotReader.java diff --git a/fluss-server/pom.xml b/fluss-server/pom.xml index 42495c253ce..3697ee0287a 100644 --- a/fluss-server/pom.xml +++ b/fluss-server/pom.xml @@ -94,6 +94,13 @@ test + + com.ververica + frocksdbjni + ${frocksdb.version} + test + + org.apache.fluss fluss-test-utils @@ -168,4 +175,4 @@ - \ No newline at end of file + diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java index c410d1bb337..a5a99a2f69e 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java @@ -29,6 +29,7 @@ import org.fluss.rocksdb.BlockBasedTableConfig; import org.fluss.rocksdb.BloomFilter; import org.fluss.rocksdb.Cache; +import org.fluss.rocksdb.ChecksumType; import org.fluss.rocksdb.ColumnFamilyOptions; import org.fluss.rocksdb.CompactionStyle; import org.fluss.rocksdb.CompressionType; @@ -73,6 +74,8 @@ public class RocksDBResourceContainer implements AutoCloseable { // the filename length limit is 255 on most operating systems private static final int INSTANCE_PATH_LENGTH_LIMIT = 255 - "_LOG".length(); + private static final int FROCKSDB_COMPATIBLE_FORMAT_VERSION = 5; + @Nullable private final File instanceRocksDBPath; /** The configurations from file. */ @@ -352,6 +355,11 @@ private ColumnFamilyOptions setColumnFamilyOptionsFromConfigurableOptions( } } + // Keep snapshots readable by FRocksDB 6.20.3. + blockBasedTableConfig + .setFormatVersion(FROCKSDB_COMPATIBLE_FORMAT_VERSION) + .setChecksumType(ChecksumType.kCRC32c); + blockBasedTableConfig.setBlockSize( internalGetOption(ConfigOptions.KV_BLOCK_SIZE).getBytes()); diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/FrocksDBSnapshotReader.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/FrocksDBSnapshotReader.java new file mode 100644 index 00000000000..3a791081b37 --- /dev/null +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/FrocksDBSnapshotReader.java @@ -0,0 +1,53 @@ +/* + * 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.server.kv.snapshot; + +import org.rocksdb.Options; +import org.rocksdb.RocksDB; + +import java.nio.charset.StandardCharsets; +import java.util.Arrays; + +/** Reads a RocksDB snapshot in an isolated process using the legacy FRocksDB JNI dependency. */ +public final class FrocksDBSnapshotReader { + + private FrocksDBSnapshotReader() {} + + /** Opens a snapshot and verifies the key/value pairs supplied after the database path. */ + public static void main(String[] args) throws Exception { + if (args.length < 3 || args.length % 2 == 0) { + throw new IllegalArgumentException( + "Expected a database path followed by one or more key/value pairs."); + } + + RocksDB.loadLibrary(); + try (Options options = new Options(); + RocksDB rocksDB = RocksDB.openReadOnly(options, args[0])) { + for (int i = 1; i < args.length; i += 2) { + byte[] actualValue = rocksDB.get(args[i].getBytes(StandardCharsets.UTF_8)); + byte[] expectedValue = args[i + 1].getBytes(StandardCharsets.UTF_8); + if (!Arrays.equals(actualValue, expectedValue)) { + throw new AssertionError( + String.format( + "Unexpected value for key %s: expected %s but was %s", + args[i], args[i + 1], Arrays.toString(actualValue))); + } + } + } + } +} diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java index 0d9b8b04a51..d3d2cc07a35 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java @@ -23,8 +23,10 @@ import org.apache.fluss.metrics.ThreadSafeSimpleCounter; import org.apache.fluss.server.kv.rocksdb.RocksDBExtension; import org.apache.fluss.server.kv.rocksdb.RocksDBKv; +import org.apache.fluss.server.kv.rocksdb.RocksDBKvBuilder; import org.apache.fluss.server.testutils.KvTestUtils; import org.apache.fluss.server.utils.ResourceGuard; +import org.apache.fluss.server.utils.TestProcessBuilder; import org.apache.fluss.utils.CloseableRegistry; import org.apache.fluss.utils.FlussPaths; @@ -35,6 +37,7 @@ import org.junit.jupiter.api.extension.RegisterExtension; import org.junit.jupiter.api.io.TempDir; +import java.nio.charset.StandardCharsets; import java.nio.file.Path; import java.util.Collection; import java.util.HashMap; @@ -42,6 +45,7 @@ import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.stream.Stream; import static org.apache.fluss.server.testutils.KvTestUtils.checkSnapshotIncrementWithNewlyFiles; @@ -155,6 +159,69 @@ void testIncrementalSnapshot(@TempDir Path snapshotBaseDir, @TempDir Path snapsh } } + @Test + void testSnapshotCanBeReadByFrocksDB( + @TempDir Path snapshotBaseDir, @TempDir Path snapshotDownDir) throws Exception { + FsPath testingTabletDir = FsPath.fromLocalFile(snapshotBaseDir.toFile()); + SnapshotLocation snapshotLocation = + new SnapshotLocation( + LocalFileSystem.getSharedInstance(), + FlussPaths.remoteKvSnapshotDir(testingTabletDir, 1L), + FlussPaths.remoteKvSharedDir(testingTabletDir), + 1024); + + try (CloseableRegistry closeableRegistry = new CloseableRegistry(); + RocksIncrementalSnapshot incrementalSnapshot = createIncrementalSnapshot()) { + RocksDB rocksDB = rocksDBExtension.getRocksDb(); + rocksDB.put( + "key1".getBytes(StandardCharsets.UTF_8), + "val1".getBytes(StandardCharsets.UTF_8)); + rocksDB.put( + "key2".getBytes(StandardCharsets.UTF_8), + "val2".getBytes(StandardCharsets.UTF_8)); + + KvSnapshotHandle snapshotHandle = + snapshot(1L, incrementalSnapshot, snapshotLocation, closeableRegistry); + incrementalSnapshot.notifySnapshotComplete(1L); + + Path restoredDbPath = snapshotDownDir.resolve(RocksDBKvBuilder.DB_INSTANCE_DIR_STRING); + KvSnapshotDataDownloader snapshotDataDownloader = + new KvSnapshotDataDownloader(dataTransferThreadPool); + snapshotDataDownloader.transferAllDataToDirectory( + new KvSnapshotDownloadSpec(snapshotHandle, restoredDbPath), closeableRegistry); + + TestProcessBuilder.TestProcess frocksDBReader = null; + try { + frocksDBReader = + new TestProcessBuilder(FrocksDBSnapshotReader.class.getName()) + .addMainClassArg(restoredDbPath.toString()) + .addMainClassArg("key1") + .addMainClassArg("val1") + .addMainClassArg("key2") + .addMainClassArg("val2") + .start(); + + boolean exited = frocksDBReader.getProcess().waitFor(1, TimeUnit.MINUTES); + assertThat(exited) + .describedAs( + "FRocksDB reader process output: %s", processOutput(frocksDBReader)) + .isTrue(); + assertThat(frocksDBReader.getProcess().exitValue()) + .describedAs( + "FRocksDB reader process output: %s", processOutput(frocksDBReader)) + .isZero(); + } finally { + if (frocksDBReader != null && frocksDBReader.getProcess().isAlive()) { + frocksDBReader.destroy(); + } + } + } + } + + private String processOutput(TestProcessBuilder.TestProcess process) { + return process.getProcessOutput().toString() + process.getErrorOutput().toString(); + } + private void verifyShareFileEqual( KvSnapshotHandle kvSnapshotHandle1, KvSnapshotHandle kvSnapshotHandle2) { List handles1 = kvSnapshotHandle1.getSharedKvFileHandles(); diff --git a/pom.xml b/pom.xml index d6367b5a0fa..f79914c169b 100644 --- a/pom.xml +++ b/pom.xml @@ -107,6 +107,7 @@ 3.1.0 3.4.3 + 6.20.3-ververica-2.0 11.8.1-fluss-2 1.7.36 2.25.4 From 927a5334329016d0980a619c44107e5a5b210107 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=99=BD=E9=B5=BA?= Date: Fri, 28 Aug 2026 17:41:32 +0800 Subject: [PATCH 3/6] [server] Avoid resetting RocksDB native loader state --- .../server/kv/rocksdb/RocksDBKvBuilder.java | 18 ------------------ .../kv/rocksdb/RocksDBKvBuilderTest.java | 9 --------- 2 files changed, 27 deletions(-) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java index aadab2a54b5..0c615d1b784 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java @@ -35,7 +35,6 @@ import java.io.File; import java.io.IOException; -import java.lang.reflect.Field; import java.util.UUID; import java.util.function.Supplier; @@ -244,15 +243,6 @@ static void ensureRocksDBIsLoaded( lastException = t; LOG.debug("RocksDB JNI library loading attempt {} failed", attempt, t); - // try to force RocksDB to attempt reloading the library - try { - resetRocksDBLoadedFlag(); - } catch (Throwable tt) { - LOG.debug( - "Failed to reset 'initialized' flag in RocksDB native code loader", - tt); - } - FileUtils.deleteDirectoryQuietly(rocksLibFolder); } } @@ -262,14 +252,6 @@ static void ensureRocksDBIsLoaded( } } - @VisibleForTesting - static void resetRocksDBLoadedFlag() throws Exception { - final Field initField = - org.fluss.rocksdb.NativeLibraryLoader.class.getDeclaredField("initialized"); - initField.setAccessible(true); - initField.setBoolean(null, false); - } - @VisibleForTesting static void resetRocksDbInitialized() { rocksDbInitialized = false; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java index 23be6d7c8e0..34002c7d1cc 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java @@ -39,15 +39,6 @@ /** Test for {@link org.apache.fluss.server.kv.rocksdb.RocksDBKvBuilder} . */ class RocksDBKvBuilderTest { - /** - * This test checks that the RocksDB native code loader still responds to resetting the init - * flag. - */ - @Test - void testResetInitFlag() throws Exception { - RocksDBKvBuilder.resetRocksDBLoadedFlag(); - } - @Test void testTempLibFolderDeletedOnFail(@TempDir Path tempDir) { RocksDBKvBuilder.resetRocksDbInitialized(); From 4fab00982aba74cdf2fc4a0c80c51678d7fa2742 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=99=BD=E9=B5=BA?= Date: Sat, 29 Aug 2026 09:27:17 +0800 Subject: [PATCH 4/6] [build] Update RocksDB JNI namespace Upgrade fluss-rocksdbjni to 11.8.1-fluss-3 and migrate imports to the relocated io.github.fluss_contrib.rocksdb namespace. --- .../scanner/batch/SnapshotFilesReader.java | 12 +++---- .../src/main/resources/META-INF/NOTICE | 2 +- .../apache/fluss/rocksdb/RocksDBHandle.java | 10 +++--- .../fluss/rocksdb/RocksDBOperationUtils.java | 12 +++---- .../fluss/rocksdb/RocksIteratorWrapper.java | 8 ++--- .../org/apache/fluss/server/kv/KvManager.java | 12 +++---- .../org/apache/fluss/server/kv/KvTablet.java | 16 ++++----- .../kv/RowTtlCompactionFilterFactory.java | 4 +-- .../fluss/server/kv/rocksdb/RocksDBKv.java | 22 ++++++------ .../server/kv/rocksdb/RocksDBKvBuilder.java | 12 +++---- .../kv/rocksdb/RocksDBResourceContainer.java | 34 +++++++++---------- .../server/kv/rocksdb/RocksDBStatistics.java | 20 +++++------ .../kv/rocksdb/RocksDBWriteBatchWrapper.java | 10 +++--- .../fluss/server/kv/scan/ScannerContext.java | 8 ++--- .../kv/snapshot/RocksIncrementalSnapshot.java | 4 +-- .../src/main/resources/META-INF/NOTICE | 2 +- .../apache/fluss/server/kv/KvTabletTest.java | 6 ++-- .../server/kv/RowTtlCompactionFilterTest.java | 8 ++--- .../server/kv/rocksdb/RocksDBExtension.java | 2 +- .../kv/rocksdb/RocksDBKvBuilderTest.java | 2 +- .../server/kv/rocksdb/RocksDBKvTest.java | 8 ++--- .../server/kv/rocksdb/RocksDBKvTestUtils.java | 2 +- .../rocksdb/RocksDBOperationsUtilsTest.java | 10 +++--- .../rocksdb/RocksDBResourceContainerTest.java | 30 ++++++++-------- .../snapshot/KvTabletSnapshotTargetTest.java | 2 +- .../RocksIncrementalSnapshotTest.java | 2 +- .../fluss/server/testutils/KvTestUtils.java | 4 +-- pom.xml | 2 +- 28 files changed, 133 insertions(+), 133 deletions(-) diff --git a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SnapshotFilesReader.java b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SnapshotFilesReader.java index 65f823bf87f..b3832a03a10 100644 --- a/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SnapshotFilesReader.java +++ b/fluss-client/src/main/java/org/apache/fluss/client/table/scanner/batch/SnapshotFilesReader.java @@ -34,12 +34,12 @@ import org.apache.fluss.utils.IOUtils; import org.apache.fluss.utils.SchemaUtil; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.DBOptions; -import org.fluss.rocksdb.ReadOptions; -import org.fluss.rocksdb.RocksDB; -import org.fluss.rocksdb.RocksIterator; -import org.fluss.rocksdb.Snapshot; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.DBOptions; +import io.github.fluss_contrib.rocksdb.ReadOptions; +import io.github.fluss_contrib.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksIterator; +import io.github.fluss_contrib.rocksdb.Snapshot; import javax.annotation.Nullable; import javax.annotation.concurrent.NotThreadSafe; diff --git a/fluss-client/src/main/resources/META-INF/NOTICE b/fluss-client/src/main/resources/META-INF/NOTICE index 6f1fe5166ff..b2379a9a60b 100644 --- a/fluss-client/src/main/resources/META-INF/NOTICE +++ b/fluss-client/src/main/resources/META-INF/NOTICE @@ -7,7 +7,7 @@ The Apache Software Foundation (http://www.apache.org/). This project bundles the following dependencies under the Apache Software License 2.0 (http://www.apache.org/licenses/LICENSE-2.0.txt) - com.google.code.findbugs:jsr305:1.3.9 -- io.github.fluss-contrib:fluss-rocksdbjni:11.8.1-fluss-2 +- io.github.fluss-contrib:fluss-rocksdbjni:11.8.1-fluss-3 - org.apache.commons:commons-lang3:3.18.0 - org.apache.commons:commons-math3:3.6.1 - at.yawk.lz4:lz4-java:1.10.2 diff --git a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBHandle.java b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBHandle.java index 6b81747923d..f4cce9654c8 100644 --- a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBHandle.java +++ b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBHandle.java @@ -19,11 +19,11 @@ import org.apache.fluss.utils.IOUtils; -import org.fluss.rocksdb.ColumnFamilyDescriptor; -import org.fluss.rocksdb.ColumnFamilyHandle; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.DBOptions; -import org.fluss.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.ColumnFamilyDescriptor; +import io.github.fluss_contrib.rocksdb.ColumnFamilyHandle; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.DBOptions; +import io.github.fluss_contrib.rocksdb.RocksDB; import java.io.File; import java.io.IOException; diff --git a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBOperationUtils.java b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBOperationUtils.java index 4060cfe37c5..ff02770a8e5 100644 --- a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBOperationUtils.java +++ b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksDBOperationUtils.java @@ -20,12 +20,12 @@ import org.apache.fluss.utils.IOUtils; import org.apache.fluss.utils.OperatingSystem; -import org.fluss.rocksdb.ColumnFamilyDescriptor; -import org.fluss.rocksdb.ColumnFamilyHandle; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.DBOptions; -import org.fluss.rocksdb.RocksDB; -import org.fluss.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.ColumnFamilyDescriptor; +import io.github.fluss_contrib.rocksdb.ColumnFamilyHandle; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.DBOptions; +import io.github.fluss_contrib.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksDBException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksIteratorWrapper.java b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksIteratorWrapper.java index 849964ffb31..b4e639675da 100644 --- a/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksIteratorWrapper.java +++ b/fluss-common/src/main/java/org/apache/fluss/rocksdb/RocksIteratorWrapper.java @@ -17,10 +17,10 @@ package org.apache.fluss.rocksdb; -import org.fluss.rocksdb.RocksDBException; -import org.fluss.rocksdb.RocksIterator; -import org.fluss.rocksdb.RocksIteratorInterface; -import org.fluss.rocksdb.Snapshot; +import io.github.fluss_contrib.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.RocksIterator; +import io.github.fluss_contrib.rocksdb.RocksIteratorInterface; +import io.github.fluss_contrib.rocksdb.Snapshot; import javax.annotation.Nonnull; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java index f3a17bf6a0d..67178b2841e 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java @@ -54,12 +54,12 @@ import org.apache.fluss.utils.function.SupplierWithException; import org.apache.fluss.utils.types.Tuple2; -import org.fluss.rocksdb.Cache; -import org.fluss.rocksdb.LRUCache; -import org.fluss.rocksdb.RateLimiter; -import org.fluss.rocksdb.RateLimiterMode; -import org.fluss.rocksdb.RocksDB; -import org.fluss.rocksdb.WriteBufferManager; +import io.github.fluss_contrib.rocksdb.Cache; +import io.github.fluss_contrib.rocksdb.LRUCache; +import io.github.fluss_contrib.rocksdb.RateLimiter; +import io.github.fluss_contrib.rocksdb.RateLimiterMode; +import io.github.fluss_contrib.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.WriteBufferManager; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java index d39df7a7787..6f43629b6ae 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java @@ -70,14 +70,14 @@ import org.apache.fluss.utils.clock.Clock; import org.apache.fluss.utils.clock.SystemClock; -import org.fluss.rocksdb.AbstractCompactionFilter; -import org.fluss.rocksdb.AbstractCompactionFilterFactory; -import org.fluss.rocksdb.Cache; -import org.fluss.rocksdb.RateLimiter; -import org.fluss.rocksdb.ReadOptions; -import org.fluss.rocksdb.RocksIterator; -import org.fluss.rocksdb.Snapshot; -import org.fluss.rocksdb.WriteBufferManager; +import io.github.fluss_contrib.rocksdb.AbstractCompactionFilter; +import io.github.fluss_contrib.rocksdb.AbstractCompactionFilterFactory; +import io.github.fluss_contrib.rocksdb.Cache; +import io.github.fluss_contrib.rocksdb.RateLimiter; +import io.github.fluss_contrib.rocksdb.ReadOptions; +import io.github.fluss_contrib.rocksdb.RocksIterator; +import io.github.fluss_contrib.rocksdb.Snapshot; +import io.github.fluss_contrib.rocksdb.WriteBufferManager; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/RowTtlCompactionFilterFactory.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/RowTtlCompactionFilterFactory.java index 2c1698cf33a..a7af7720151 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/RowTtlCompactionFilterFactory.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/RowTtlCompactionFilterFactory.java @@ -21,8 +21,8 @@ import org.apache.fluss.server.utils.RowTtlUtils; import org.apache.fluss.utils.clock.Clock; -import org.fluss.rocksdb.FlussTtlCompactionFilter; -import org.fluss.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.FlussTtlCompactionFilter; +import io.github.fluss_contrib.rocksdb.RocksDB; import java.time.Duration; import java.util.function.LongSupplier; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java index fe3537da139..106daf57b2d 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKv.java @@ -27,16 +27,16 @@ import org.apache.fluss.utils.BytesUtils; import org.apache.fluss.utils.IOUtils; -import org.fluss.rocksdb.Cache; -import org.fluss.rocksdb.ColumnFamilyHandle; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.MutableDBOptions; -import org.fluss.rocksdb.ReadOptions; -import org.fluss.rocksdb.RocksDB; -import org.fluss.rocksdb.RocksDBException; -import org.fluss.rocksdb.RocksIterator; -import org.fluss.rocksdb.Statistics; -import org.fluss.rocksdb.WriteOptions; +import io.github.fluss_contrib.rocksdb.Cache; +import io.github.fluss_contrib.rocksdb.ColumnFamilyHandle; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.MutableDBOptions; +import io.github.fluss_contrib.rocksdb.ReadOptions; +import io.github.fluss_contrib.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.RocksIterator; +import io.github.fluss_contrib.rocksdb.Statistics; +import io.github.fluss_contrib.rocksdb.WriteOptions; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -49,7 +49,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; -/** A wrapper for the operation of {@link org.fluss.rocksdb.RocksDB}. */ +/** A wrapper for the operation of {@link io.github.fluss_contrib.rocksdb.RocksDB}. */ public class RocksDBKv implements AutoCloseable { private static final Logger LOG = LoggerFactory.getLogger(RocksDBKv.class); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java index 0c615d1b784..12f31474389 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilder.java @@ -24,12 +24,12 @@ import org.apache.fluss.utils.FileUtils; import org.apache.fluss.utils.IOUtils; -import org.fluss.rocksdb.AbstractCompactionFilter; -import org.fluss.rocksdb.AbstractCompactionFilterFactory; -import org.fluss.rocksdb.ColumnFamilyHandle; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.NativeLibraryLoader; -import org.fluss.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.AbstractCompactionFilter; +import io.github.fluss_contrib.rocksdb.AbstractCompactionFilterFactory; +import io.github.fluss_contrib.rocksdb.ColumnFamilyHandle; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.NativeLibraryLoader; +import io.github.fluss_contrib.rocksdb.RocksDB; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java index a5a99a2f69e..46bc190a13c 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainer.java @@ -26,23 +26,23 @@ import org.apache.fluss.utils.FileUtils; import org.apache.fluss.utils.IOUtils; -import org.fluss.rocksdb.BlockBasedTableConfig; -import org.fluss.rocksdb.BloomFilter; -import org.fluss.rocksdb.Cache; -import org.fluss.rocksdb.ChecksumType; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.CompactionStyle; -import org.fluss.rocksdb.CompressionType; -import org.fluss.rocksdb.DBOptions; -import org.fluss.rocksdb.InfoLogLevel; -import org.fluss.rocksdb.LRUCache; -import org.fluss.rocksdb.PlainTableConfig; -import org.fluss.rocksdb.RateLimiter; -import org.fluss.rocksdb.ReadOptions; -import org.fluss.rocksdb.Statistics; -import org.fluss.rocksdb.TableFormatConfig; -import org.fluss.rocksdb.WriteBufferManager; -import org.fluss.rocksdb.WriteOptions; +import io.github.fluss_contrib.rocksdb.BlockBasedTableConfig; +import io.github.fluss_contrib.rocksdb.BloomFilter; +import io.github.fluss_contrib.rocksdb.Cache; +import io.github.fluss_contrib.rocksdb.ChecksumType; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.CompactionStyle; +import io.github.fluss_contrib.rocksdb.CompressionType; +import io.github.fluss_contrib.rocksdb.DBOptions; +import io.github.fluss_contrib.rocksdb.InfoLogLevel; +import io.github.fluss_contrib.rocksdb.LRUCache; +import io.github.fluss_contrib.rocksdb.PlainTableConfig; +import io.github.fluss_contrib.rocksdb.RateLimiter; +import io.github.fluss_contrib.rocksdb.ReadOptions; +import io.github.fluss_contrib.rocksdb.Statistics; +import io.github.fluss_contrib.rocksdb.TableFormatConfig; +import io.github.fluss_contrib.rocksdb.WriteBufferManager; +import io.github.fluss_contrib.rocksdb.WriteOptions; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBStatistics.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBStatistics.java index dc1407a172d..21e8ba23359 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBStatistics.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBStatistics.java @@ -19,16 +19,16 @@ import org.apache.fluss.server.utils.ResourceGuard; -import org.fluss.rocksdb.Cache; -import org.fluss.rocksdb.ColumnFamilyHandle; -import org.fluss.rocksdb.HistogramData; -import org.fluss.rocksdb.HistogramType; -import org.fluss.rocksdb.MemoryUsageType; -import org.fluss.rocksdb.MemoryUtil; -import org.fluss.rocksdb.RocksDB; -import org.fluss.rocksdb.RocksDBException; -import org.fluss.rocksdb.Statistics; -import org.fluss.rocksdb.TickerType; +import io.github.fluss_contrib.rocksdb.Cache; +import io.github.fluss_contrib.rocksdb.ColumnFamilyHandle; +import io.github.fluss_contrib.rocksdb.HistogramData; +import io.github.fluss_contrib.rocksdb.HistogramType; +import io.github.fluss_contrib.rocksdb.MemoryUsageType; +import io.github.fluss_contrib.rocksdb.MemoryUtil; +import io.github.fluss_contrib.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.Statistics; +import io.github.fluss_contrib.rocksdb.TickerType; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBWriteBatchWrapper.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBWriteBatchWrapper.java index cf34a2c1831..ea4d0de7f90 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBWriteBatchWrapper.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rocksdb/RocksDBWriteBatchWrapper.java @@ -24,11 +24,11 @@ import org.apache.fluss.server.kv.KvBatchWriter; import org.apache.fluss.utils.IOUtils; -import org.fluss.rocksdb.RocksDB; -import org.fluss.rocksdb.RocksDBException; -import org.fluss.rocksdb.Status; -import org.fluss.rocksdb.WriteBatch; -import org.fluss.rocksdb.WriteOptions; +import io.github.fluss_contrib.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.Status; +import io.github.fluss_contrib.rocksdb.WriteBatch; +import io.github.fluss_contrib.rocksdb.WriteOptions; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/scan/ScannerContext.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/scan/ScannerContext.java index 15201cf28f9..4b1befc7fd8 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/scan/ScannerContext.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/scan/ScannerContext.java @@ -24,10 +24,10 @@ import org.apache.fluss.server.utils.ResourceGuard; import org.apache.fluss.utils.IOUtils; -import org.fluss.rocksdb.ReadOptions; -import org.fluss.rocksdb.RocksDBException; -import org.fluss.rocksdb.RocksIterator; -import org.fluss.rocksdb.Snapshot; +import io.github.fluss_contrib.rocksdb.ReadOptions; +import io.github.fluss_contrib.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.RocksIterator; +import io.github.fluss_contrib.rocksdb.Snapshot; import javax.annotation.concurrent.NotThreadSafe; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshot.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshot.java index 56e28fdf30f..af43b9ed833 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshot.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshot.java @@ -23,8 +23,8 @@ import org.apache.fluss.utils.ExceptionUtils; import org.apache.fluss.utils.FileUtils; -import org.fluss.rocksdb.Checkpoint; -import org.fluss.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.Checkpoint; +import io.github.fluss_contrib.rocksdb.RocksDB; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/fluss-server/src/main/resources/META-INF/NOTICE b/fluss-server/src/main/resources/META-INF/NOTICE index a6c20fe7b08..a29069ac81e 100644 --- a/fluss-server/src/main/resources/META-INF/NOTICE +++ b/fluss-server/src/main/resources/META-INF/NOTICE @@ -9,7 +9,7 @@ This project bundles the following dependencies under the Apache Software Licens - com.github.ben-manes.caffeine:caffeine:2.9.3 - com.google.code.findbugs:jsr305:1.3.9 - com.google.errorprone:error_prone_annotations:2.10.0 -- io.github.fluss-contrib:fluss-rocksdbjni:11.8.1-fluss-2 +- io.github.fluss-contrib:fluss-rocksdbjni:11.8.1-fluss-3 - commons-cli:commons-cli:1.5.0 - org.apache.commons:commons-lang3:3.18.0 - org.apache.commons:commons-math3:3.6.1 diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvTabletTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvTabletTest.java index 8a9d2752a49..effc54110cc 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvTabletTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvTabletTest.java @@ -87,8 +87,8 @@ import org.apache.fluss.utils.clock.SystemClock; import org.apache.fluss.utils.concurrent.FlussScheduler; -import org.fluss.rocksdb.FlushOptions; -import org.fluss.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.FlushOptions; +import io.github.fluss_contrib.rocksdb.RocksDBException; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -2285,7 +2285,7 @@ void testRocksDBMetrics() throws Exception { assertThat(statistics).as("RocksDB statistics should be available").isNotNull(); // Verify statistics is properly initialized - org.fluss.rocksdb.Statistics stats = kvTablet.getRocksDBKv().getStatistics(); + io.github.fluss_contrib.rocksdb.Statistics stats = kvTablet.getRocksDBKv().getStatistics(); assertThat(stats).as("RocksDB Statistics should be enabled").isNotNull(); // All metrics should start at 0 for a fresh database diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/RowTtlCompactionFilterTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/RowTtlCompactionFilterTest.java index c5e83bd0209..01c6fa064c9 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/RowTtlCompactionFilterTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/RowTtlCompactionFilterTest.java @@ -25,10 +25,10 @@ import org.apache.fluss.row.encode.ValueEncoder; import org.apache.fluss.utils.clock.ManualClock; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.DBOptions; -import org.fluss.rocksdb.FlushOptions; -import org.fluss.rocksdb.FlussTtlCompactionFilter; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.DBOptions; +import io.github.fluss_contrib.rocksdb.FlushOptions; +import io.github.fluss_contrib.rocksdb.FlussTtlCompactionFilter; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBExtension.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBExtension.java index cf3fcf78d28..9fe4b21409b 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBExtension.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBExtension.java @@ -20,7 +20,7 @@ import org.apache.fluss.config.Configuration; import org.apache.fluss.utils.IOUtils; -import org.fluss.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksDB; import org.junit.jupiter.api.extension.AfterEachCallback; import org.junit.jupiter.api.extension.BeforeEachCallback; import org.junit.jupiter.api.extension.ExtensionContext; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java index 34002c7d1cc..8e6fb96101d 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvBuilderTest.java @@ -24,7 +24,7 @@ import org.apache.fluss.server.kv.RowTtlCompactionFilterFactory; import org.apache.fluss.utils.clock.ManualClock; -import org.fluss.rocksdb.FlussTtlCompactionFilter; +import io.github.fluss_contrib.rocksdb.FlussTtlCompactionFilter; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java index ac05209afa3..a0035d5570e 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTest.java @@ -23,9 +23,9 @@ import org.apache.fluss.metrics.util.TestHistogram; import org.apache.fluss.server.kv.KvCloseMode; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.FlushOptions; -import org.fluss.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.FlushOptions; +import io.github.fluss_contrib.rocksdb.RocksDBException; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; @@ -130,7 +130,7 @@ void testNoSlowdownWriteBatchFailsFastOnRocksDbDelay(@TempDir Path tempDir) thro } /** - * Verifies that the L0 property is actually readable on fluss-rocksdbjni 11.8.1-fluss-2. + * Verifies that the L0 property is actually readable on fluss-rocksdbjni 11.8.1-fluss-3. * *

This is a regression test: {@code getLongProperty("rocksdb.num-files-at-level0")} throws * {@code RocksDBException: NotFound} because the property is parametric (string-type), not an diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTestUtils.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTestUtils.java index 75646ec414f..62e96795d6f 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTestUtils.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBKvTestUtils.java @@ -18,7 +18,7 @@ package org.apache.fluss.server.kv.rocksdb; -import org.fluss.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.RocksDBException; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.spy; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBOperationsUtilsTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBOperationsUtilsTest.java index adbb8a515c8..2edc0586192 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBOperationsUtilsTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBOperationsUtilsTest.java @@ -19,11 +19,11 @@ import org.apache.fluss.rocksdb.RocksDBOperationUtils; -import org.fluss.rocksdb.ColumnFamilyDescriptor; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.DBOptions; -import org.fluss.rocksdb.RocksDB; -import org.fluss.rocksdb.RocksDBException; +import io.github.fluss_contrib.rocksdb.ColumnFamilyDescriptor; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.DBOptions; +import io.github.fluss_contrib.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksDBException; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java index 2dd117d736d..03a9eaea434 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rocksdb/RocksDBResourceContainerTest.java @@ -21,21 +21,21 @@ import org.apache.fluss.config.Configuration; import org.apache.fluss.server.kv.KvManager; -import org.fluss.rocksdb.BlockBasedTableConfig; -import org.fluss.rocksdb.BloomFilter; -import org.fluss.rocksdb.Cache; -import org.fluss.rocksdb.ColumnFamilyOptions; -import org.fluss.rocksdb.CompactionStyle; -import org.fluss.rocksdb.CompressionType; -import org.fluss.rocksdb.DBOptions; -import org.fluss.rocksdb.FlushOptions; -import org.fluss.rocksdb.InfoLogLevel; -import org.fluss.rocksdb.LRUCache; -import org.fluss.rocksdb.RateLimiter; -import org.fluss.rocksdb.ReadOptions; -import org.fluss.rocksdb.WriteBufferManager; -import org.fluss.rocksdb.WriteOptions; -import org.fluss.rocksdb.util.SizeUnit; +import io.github.fluss_contrib.rocksdb.BlockBasedTableConfig; +import io.github.fluss_contrib.rocksdb.BloomFilter; +import io.github.fluss_contrib.rocksdb.Cache; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.CompactionStyle; +import io.github.fluss_contrib.rocksdb.CompressionType; +import io.github.fluss_contrib.rocksdb.DBOptions; +import io.github.fluss_contrib.rocksdb.FlushOptions; +import io.github.fluss_contrib.rocksdb.InfoLogLevel; +import io.github.fluss_contrib.rocksdb.LRUCache; +import io.github.fluss_contrib.rocksdb.RateLimiter; +import io.github.fluss_contrib.rocksdb.ReadOptions; +import io.github.fluss_contrib.rocksdb.WriteBufferManager; +import io.github.fluss_contrib.rocksdb.WriteOptions; +import io.github.fluss_contrib.rocksdb.util.SizeUnit; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/KvTabletSnapshotTargetTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/KvTabletSnapshotTargetTest.java index 3f1124488a8..b7ceda3dfbb 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/KvTabletSnapshotTargetTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/KvTabletSnapshotTargetTest.java @@ -42,7 +42,7 @@ import org.apache.fluss.utils.FlussPaths; import org.apache.fluss.utils.concurrent.Executors; -import org.fluss.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksDB; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeEach; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java index d3d2cc07a35..d5b1a7593a6 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/RocksIncrementalSnapshotTest.java @@ -30,7 +30,7 @@ import org.apache.fluss.utils.CloseableRegistry; import org.apache.fluss.utils.FlussPaths; -import org.fluss.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksDB; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/testutils/KvTestUtils.java b/fluss-server/src/test/java/org/apache/fluss/server/testutils/KvTestUtils.java index e735e681d65..34e74356ab4 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/testutils/KvTestUtils.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/testutils/KvTestUtils.java @@ -39,8 +39,8 @@ import org.apache.fluss.utils.FileUtils; import org.apache.fluss.utils.types.Tuple2; -import org.fluss.rocksdb.RocksDB; -import org.fluss.rocksdb.RocksIterator; +import io.github.fluss_contrib.rocksdb.RocksDB; +import io.github.fluss_contrib.rocksdb.RocksIterator; import javax.annotation.Nullable; diff --git a/pom.xml b/pom.xml index f79914c169b..5f9466c8eba 100644 --- a/pom.xml +++ b/pom.xml @@ -108,7 +108,7 @@ 3.4.3 6.20.3-ververica-2.0 - 11.8.1-fluss-2 + 11.8.1-fluss-3 1.7.36 2.25.4 2.3.1 From 0953ba5afde137631e4a484fabf97ac12b6e0791 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=99=BD=E9=B5=BA?= Date: Sun, 27 Sep 2026 08:46:38 +0800 Subject: [PATCH 5/6] [test] Align new tests with relocated RocksDB JNI --- .../server/kv/HistoricalKvCompactionFilterTest.java | 10 +++++----- .../historical/HistoricalPartitionManagerTest.java | 2 +- .../server/tablet/TabletServerShutdownITCase.java | 2 +- 3 files changed, 7 insertions(+), 7 deletions(-) diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/HistoricalKvCompactionFilterTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/HistoricalKvCompactionFilterTest.java index 0b076974128..09bf435bdf8 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/HistoricalKvCompactionFilterTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/HistoricalKvCompactionFilterTest.java @@ -24,13 +24,13 @@ import org.apache.fluss.row.encode.ValueEncoder; import org.apache.fluss.server.kv.historical.HistoricalKvTombstone; +import io.github.fluss_contrib.rocksdb.ColumnFamilyOptions; +import io.github.fluss_contrib.rocksdb.DBOptions; +import io.github.fluss_contrib.rocksdb.FlushOptions; +import io.github.fluss_contrib.rocksdb.FlussTtlCompactionFilter; import org.junit.jupiter.api.io.TempDir; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; -import org.rocksdb.ColumnFamilyOptions; -import org.rocksdb.DBOptions; -import org.rocksdb.FlinkCompactionFilter; -import org.rocksdb.FlushOptions; import java.nio.charset.StandardCharsets; import java.nio.file.Path; @@ -59,7 +59,7 @@ void testRemovesOnlyOffsetsBeyondRetentionDistance(long retentionOffsetDistance) byte[] tombstoneAfterKey = bytes("tombstone-after"); BinaryRow row = compactedRow(DATA1_ROW_TYPE, new Object[] {1, "a"}); - try (FlinkCompactionFilter.FlinkCompactionFilterFactory filterFactory = + try (FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory filterFactory = RowTtlCompactionFilterFactory.create( KvValueLayout.TAGGED, retentionOffsetDistance, diff --git a/fluss-server/src/test/java/org/apache/fluss/server/replica/historical/HistoricalPartitionManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/replica/historical/HistoricalPartitionManagerTest.java index 73d41737d8b..553e7382b78 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/replica/historical/HistoricalPartitionManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/replica/historical/HistoricalPartitionManagerTest.java @@ -90,10 +90,10 @@ import com.github.benmanes.caffeine.cache.Scheduler; import com.github.benmanes.caffeine.cache.Ticker; +import io.github.fluss_contrib.rocksdb.FlushOptions; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; -import org.rocksdb.FlushOptions; import javax.annotation.Nullable; diff --git a/fluss-server/src/test/java/org/apache/fluss/server/tablet/TabletServerShutdownITCase.java b/fluss-server/src/test/java/org/apache/fluss/server/tablet/TabletServerShutdownITCase.java index d1921f7df4e..5998e489f17 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/tablet/TabletServerShutdownITCase.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/tablet/TabletServerShutdownITCase.java @@ -36,12 +36,12 @@ import org.apache.fluss.utils.FlussPaths; import org.apache.fluss.utils.types.Tuple2; +import io.github.fluss_contrib.rocksdb.FlushOptions; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.CsvSource; import org.junit.jupiter.params.provider.ValueSource; -import org.rocksdb.FlushOptions; import java.io.File; import java.nio.file.Files; From 73777925c4f6b2c20cdbb7081d90e35b220a2025 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=99=BD=E9=B5=BA?= Date: Sun, 27 Sep 2026 09:23:59 +0800 Subject: [PATCH 6/6] [flink] Isolate bounded snapshot split test --- .../source/enumerator/FlinkSourceEnumeratorTest.java | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/enumerator/FlinkSourceEnumeratorTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/enumerator/FlinkSourceEnumeratorTest.java index 43764059341..fe2f5dcf7cb 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/enumerator/FlinkSourceEnumeratorTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/source/enumerator/FlinkSourceEnumeratorTest.java @@ -214,13 +214,17 @@ void testBoundedPkTableEmitsKvBatchSplits() throws Throwable { @Test void testBoundedPkTableEmitsSnapshotSplitsByDefault() throws Throwable { - createTable(DEFAULT_TABLE_PATH, DEFAULT_PK_TABLE_DESCRIPTOR); + TablePath tablePath = TablePath.of(DEFAULT_DB, "bounded-pk-snapshot-default"); + long tableId = createTable(tablePath, DEFAULT_PK_TABLE_DESCRIPTOR); + for (int bucket = 0; bucket < DEFAULT_BUCKET_NUM; bucket++) { + FLUSS_CLUSTER_EXTENSION.waitUntilAllReplicaReady(new TableBucket(tableId, bucket)); + } int numSubtasks = DEFAULT_BUCKET_NUM; try (MockSplitEnumeratorContext context = new MockSplitEnumeratorContext<>(numSubtasks)) { FlinkSourceEnumerator enumerator = new FlinkSourceEnumerator( - DEFAULT_TABLE_PATH, + tablePath, flussConf, true, false,