Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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 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;
Expand Down
4 changes: 2 additions & 2 deletions fluss-client/src/main/resources/META-INF/NOTICE
Original file line number Diff line number Diff line change
Expand Up @@ -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-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
Expand All @@ -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
- com.github.luben:zstd-jni:1.5.7-6
4 changes: 2 additions & 2 deletions fluss-common/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -98,8 +98,8 @@
the rocksdb should be provided as a kv plugin to used by client & server.
-->
<dependency>
<groupId>com.ververica</groupId>
<artifactId>frocksdbjni</artifactId>
<groupId>io.github.fluss-contrib</groupId>
<artifactId>fluss-rocksdbjni</artifactId>
</dependency>

<!-- test dependencies -->
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,9 +17,10 @@

package org.apache.fluss.rocksdb;

import org.rocksdb.RocksDBException;
import org.rocksdb.RocksIterator;
import org.rocksdb.RocksIteratorInterface;
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;

Expand Down Expand Up @@ -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();
}
Expand Down
9 changes: 8 additions & 1 deletion fluss-server/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,13 @@
<scope>test</scope>
</dependency>

<dependency>
<groupId>com.ververica</groupId>
<artifactId>frocksdbjni</artifactId>
<version>${frocksdb.version}</version>
<scope>test</scope>
</dependency>

<dependency>
<groupId>org.apache.fluss</groupId>
<artifactId>fluss-test-utils</artifactId>
Expand Down Expand Up @@ -168,4 +175,4 @@
</plugins>
</build>

</project>
</project>
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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.");
Expand All @@ -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));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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);
}
}
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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. */
Expand Down Expand Up @@ -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());

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Loading
Loading