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/KvManager.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvManager.java
index 9f00156e244..faa2b2afbcf 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
@@ -53,9 +53,9 @@
import org.apache.fluss.utils.clock.SystemClock;
import org.apache.fluss.utils.types.Tuple2;
-import org.rocksdb.RateLimiter;
-import org.rocksdb.RateLimiterMode;
-import org.rocksdb.RocksDB;
+import io.github.fluss_contrib.rocksdb.RateLimiter;
+import io.github.fluss_contrib.rocksdb.RateLimiterMode;
+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/KvTablet.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java
index 6118bb40fb3..ff88d816107 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
@@ -67,12 +67,12 @@
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.RateLimiter;
-import org.rocksdb.ReadOptions;
-import org.rocksdb.RocksIterator;
-import org.rocksdb.Snapshot;
+import io.github.fluss_contrib.rocksdb.AbstractCompactionFilter;
+import io.github.fluss_contrib.rocksdb.AbstractCompactionFilterFactory;
+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 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 8cb01d0c1f2..1fcd4bc8acb 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
@@ -22,8 +22,8 @@
import org.apache.fluss.server.utils.RowTtlUtils;
import org.apache.fluss.utils.clock.Clock;
-import org.rocksdb.FlinkCompactionFilter;
-import org.rocksdb.RocksDB;
+import io.github.fluss_contrib.rocksdb.FlussTtlCompactionFilter;
+import io.github.fluss_contrib.rocksdb.RocksDB;
import java.time.Duration;
@@ -38,13 +38,13 @@ 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) {
return create(kvValueLayout, ttl, QUERY_TIME_AFTER_NUM_ENTRIES, clock);
}
@VisibleForTesting
- static FlinkCompactionFilter.FlinkCompactionFilterFactory create(
+ static FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory create(
KvValueLayout kvValueLayout, Duration ttl, long queryTimeAfterNumEntries, Clock clock) {
long ttlMillis = RowTtlUtils.validateAndConvertTtlDurationToMillis(ttl);
checkNotNull(kvValueLayout, "kvValueLayout must not be null.");
@@ -55,11 +55,11 @@ static FlinkCompactionFilter.FlinkCompactionFilterFactory create(
"queryTimeAfterNumEntries must be greater than zero.");
RocksDB.loadLibrary();
- FlinkCompactionFilter.FlinkCompactionFilterFactory factory =
- new FlinkCompactionFilter.FlinkCompactionFilterFactory(clock::milliseconds);
+ FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory factory =
+ new FlussTtlCompactionFilter.FlussTtlCompactionFilterFactory(clock::milliseconds);
factory.configure(
- FlinkCompactionFilter.Config.createNotList(
- FlinkCompactionFilter.StateType.Value,
+ FlussTtlCompactionFilter.Config.createNotList(
+ FlussTtlCompactionFilter.StateType.Value,
kvValueLayout.valueTagOffset(),
ttlMillis,
queryTimeAfterNumEntries));
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..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.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 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.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 b44e892829c..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,18 +24,17 @@
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 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;
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.rocksdb.NativeLibraryLoader.class.getDeclaredField("initialized");
- initField.setAccessible(true);
- initField.setBoolean(null, false);
- }
-
@VisibleForTesting
static void resetRocksDbInitialized() {
rocksDbInitialized = 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 50f348f8db4..8187dbc44e0 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,21 +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.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.WriteOptions;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -72,6 +73,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. */
@@ -300,6 +303,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/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 3a4e4e5bf83..137109d1d24 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 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 7aaeab059c9..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.rocksdb.RocksDB;
-import org.rocksdb.RocksDBException;
-import org.rocksdb.Status;
-import org.rocksdb.WriteBatch;
-import org.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 aeb49886c9c..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.rocksdb.ReadOptions;
-import org.rocksdb.RocksDBException;
-import org.rocksdb.RocksIterator;
-import org.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/KvSnapshotDataUploader.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/snapshot/KvSnapshotDataUploader.java
index ef4d3b6eee1..314ff18a3d1 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
@@ -141,6 +141,8 @@ private KvFileHandleAndLocalPath uploadLocalFileToSnapshotLocation(
outputStream.write(buffer, 0, 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 c50fe3ef06c..4d077d44ed8 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
@@ -22,8 +22,8 @@
import org.apache.fluss.utils.ExceptionUtils;
import org.apache.fluss.utils.FileUtils;
-import org.rocksdb.Checkpoint;
-import org.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 2604b80cb38..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
-- com.ververica:frocksdbjni:6.20.3-ververica-2.0
+- 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
@@ -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 a7f33b284e0..d5f232c3d6e 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 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;
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;
@@ -2278,7 +2278,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();
+ 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 762b3756c2e..c01b6dfae88 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 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;
-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),
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..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,10 +20,10 @@
import org.apache.fluss.config.Configuration;
import org.apache.fluss.utils.IOUtils;
+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;
-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..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,9 +24,9 @@
import org.apache.fluss.server.kv.RowTtlCompactionFilterFactory;
import org.apache.fluss.utils.clock.ManualClock;
+import io.github.fluss_contrib.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;
@@ -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();
@@ -66,7 +57,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 +81,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..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,11 +23,11 @@
import org.apache.fluss.metrics.util.TestHistogram;
import org.apache.fluss.server.kv.KvCloseMode;
+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;
-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-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 74e5f466188..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.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 d2f9d8b8268..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,13 @@
import org.apache.fluss.rocksdb.RocksDBOperationUtils;
+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;
-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 533a032b7bd..5af3a10a2f1 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
@@ -20,18 +20,18 @@
import org.apache.fluss.config.ConfigOptions;
import org.apache.fluss.config.Configuration;
+import io.github.fluss_contrib.rocksdb.BlockBasedTableConfig;
+import io.github.fluss_contrib.rocksdb.BloomFilter;
+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.ReadOptions;
+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;
-import org.rocksdb.BlockBasedTableConfig;
-import org.rocksdb.BloomFilter;
-import org.rocksdb.ColumnFamilyOptions;
-import org.rocksdb.CompactionStyle;
-import org.rocksdb.CompressionType;
-import org.rocksdb.DBOptions;
-import org.rocksdb.InfoLogLevel;
-import org.rocksdb.ReadOptions;
-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/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/KvTabletSnapshotTargetTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/snapshot/KvTabletSnapshotTargetTest.java
index 29773970ea6..97ae9405c99 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
@@ -40,13 +40,13 @@
import org.apache.fluss.utils.FlussPaths;
import org.apache.fluss.utils.concurrent.Executors;
+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;
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 64398516abc..58f1f737401 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
@@ -21,18 +21,21 @@
import org.apache.fluss.fs.local.LocalFileSystem;
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;
+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;
import org.junit.jupiter.api.extension.RegisterExtension;
import org.junit.jupiter.api.io.TempDir;
-import org.rocksdb.RocksDB;
+import java.nio.charset.StandardCharsets;
import java.nio.file.Path;
import java.util.Collection;
import java.util.HashMap;
@@ -40,6 +43,7 @@
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
import static org.apache.fluss.server.testutils.KvTestUtils.checkSnapshotIncrementWithNewlyFiles;
import static org.assertj.core.api.Assertions.assertThat;
@@ -148,6 +152,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/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..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.rocksdb.RocksDB;
-import org.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 480c67018f2..ca1747aa93b 100644
--- a/pom.xml
+++ b/pom.xml
@@ -108,6 +108,7 @@
3.4.0
6.20.3-ververica-2.0
+ 11.8.1-fluss-3
1.7.36
2.25.4
2.3.1
@@ -367,9 +368,9 @@