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 @@ -258,6 +258,14 @@ public class MetricNames {
public static final String ROCKSDB_SHARED_WRITE_BUFFER_CAPACITY =
"rocksdbSharedWriteBufferCapacity";

// Server-level pre-write buffer metrics (aggregated from all KV tablets, Sum aggregation)
/**
* Estimated memory usage of the pre-write buffers across all KV tablets in this server (Sum
* aggregation).
*/
public static final String KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES =
"kvPreWriteBufferMemoryUsageBytes";

// Table-level RocksDB memory metrics (Sum aggregation)
/** Total memtable memory usage across all buckets of this table. */
public static final String ROCKSDB_MEMTABLE_MEMORY_USAGE_TOTAL =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@
import java.util.Objects;
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;

import static org.apache.fluss.utils.concurrent.LockUtils.inLock;

Expand Down Expand Up @@ -138,6 +139,9 @@ public static RateLimiter getDefaultRateLimiter() {
/** The memory segment pool to allocate memorySegment. */
private final LazyMemorySegmentPool memorySegmentPool;

/** Server-wide pre-write buffer memory usage shared by all KV tablets, in bytes. */
private final AtomicLong kvPreWriteBufferMemoryUsageBytes = new AtomicLong();

private final FsPath remoteKvDir;

private final FileSystem remoteFileSystem;
Expand Down Expand Up @@ -225,6 +229,8 @@ private KvManager(
throw e;
}
this.kvFlushScheduler = createdFlushScheduler;
tabletServerMetricGroup.setKvPreWriteBufferMemoryUsageMetrics(
kvPreWriteBufferMemoryUsageBytes::get);
if (sharedBlockCache != null) {
tabletServerMetricGroup.setSharedBlockCacheMetrics(
this::getSharedBlockCacheUsage,
Expand Down Expand Up @@ -459,6 +465,7 @@ public KvTablet getOrCreateKv(
sharedRocksDBRateLimiter,
sharedBlockCache,
sharedWriteBufferManager,
kvPreWriteBufferMemoryUsageBytes,
kvFlushScheduler,
flushCompleteListener,
autoIncrementManager,
Expand Down Expand Up @@ -582,6 +589,7 @@ public KvTablet loadKv(
sharedRocksDBRateLimiter,
sharedBlockCache,
sharedWriteBufferManager,
kvPreWriteBufferMemoryUsageBytes,
kvFlushScheduler,
flushCompleteListener,
autoIncrementManager,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.Executor;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;

Expand Down Expand Up @@ -203,6 +204,7 @@ private KvTablet(
ValueEncoder valueEncoder,
ValueDecoder valueDecoder,
@Nullable RocksDBStatistics rocksDBStatistics,
AtomicLong sharedPreWriteBufferMemoryUsageBytes,
KvFlushScheduler kvFlushScheduler,
boolean closeFlushScheduler,
@Nullable Runnable flushCompleteListener,
Expand All @@ -221,7 +223,8 @@ private KvTablet(
this.serverMetricGroup = serverMetricGroup;
this.kvFlushScheduler = kvFlushScheduler;
this.closeFlushScheduler = closeFlushScheduler;
this.kvPreWriteBuffer = new KvPreWriteBuffer(serverMetricGroup);
this.kvPreWriteBuffer =
new KvPreWriteBuffer(serverMetricGroup, sharedPreWriteBufferMemoryUsageBytes);
this.kvStateAccessor =
new KvStateAccessor(kvPreWriteBuffer, rocksDBKv, historicalPartition);
this.kvValueLayout = kvValueLayout;
Expand Down Expand Up @@ -299,6 +302,7 @@ public static KvTablet create(
sharedRateLimiter,
null,
null,
new AtomicLong(),
new KvFlushScheduler(serverConf),
true,
null,
Expand All @@ -324,6 +328,7 @@ static KvTablet create(
RateLimiter sharedRateLimiter,
@Nullable Cache sharedBlockCache,
@Nullable WriteBufferManager sharedWriteBufferManager,
AtomicLong sharedPreWriteBufferMemoryUsageBytes,
KvFlushScheduler kvFlushScheduler,
@Nullable Runnable flushCompleteListener,
AutoIncrementManager autoIncrementManager,
Expand All @@ -347,6 +352,7 @@ static KvTablet create(
sharedRateLimiter,
sharedBlockCache,
sharedWriteBufferManager,
sharedPreWriteBufferMemoryUsageBytes,
kvFlushScheduler,
false,
flushCompleteListener,
Expand All @@ -371,6 +377,7 @@ public static KvTablet create(
ChangelogImage changelogImage,
RateLimiter sharedRateLimiter,
@Nullable Cache sharedBlockCache,
AtomicLong sharedPreWriteBufferMemoryUsageBytes,
KvFlushScheduler kvFlushScheduler,
@Nullable Runnable flushCompleteListener,
AutoIncrementManager autoIncrementManager,
Expand All @@ -394,6 +401,7 @@ public static KvTablet create(
sharedRateLimiter,
sharedBlockCache,
null,
sharedPreWriteBufferMemoryUsageBytes,
kvFlushScheduler,
false,
flushCompleteListener,
Expand All @@ -419,6 +427,7 @@ private static KvTablet create(
RateLimiter sharedRateLimiter,
@Nullable Cache sharedBlockCache,
@Nullable WriteBufferManager sharedWriteBufferManager,
AtomicLong sharedPreWriteBufferMemoryUsageBytes,
KvFlushScheduler kvFlushScheduler,
boolean closeFlushScheduler,
@Nullable Runnable flushCompleteListener,
Expand Down Expand Up @@ -487,6 +496,7 @@ private static KvTablet create(
valueEncoder,
valueDecoder,
rocksDBStatistics,
sharedPreWriteBufferMemoryUsageBytes,
kvFlushScheduler,
closeFlushScheduler,
flushCompleteListener,
Expand Down Expand Up @@ -532,6 +542,7 @@ public static KvTablet create(
sharedRateLimiter,
null,
null,
new AtomicLong(),
new KvFlushScheduler(serverConf),
true,
null,
Expand Down Expand Up @@ -1346,6 +1357,9 @@ public void close(KvCloseMode closeMode) throws Exception {
// Terminal transition: closing forces IDLE regardless of the current
// state, see the FlushState state graph.
flushState = FlushState.IDLE;
// Release the remaining pre-write buffer accounting to the shared
// ledger while the local accounting values are still exact.
kvPreWriteBuffer.close();
return true;
});
if (shouldClose && closeFlushScheduler) {
Expand Down
Loading
Loading