From 972f932fd5a4293830bfa53cf57fce2c69a0582d Mon Sep 17 00:00:00 2001 From: duankaixuan <1417048384@qq.com> Date: Sat, 12 Sep 2026 00:05:59 +0800 Subject: [PATCH 1/3] [server] Add pre-write buffer memory metrics for primary key tables --- .../org/apache/fluss/metrics/MetricNames.java | 14 ++++ .../org/apache/fluss/server/kv/KvManager.java | 9 +++ .../org/apache/fluss/server/kv/KvTablet.java | 10 +++ .../server/kv/prewrite/KvPreWriteBuffer.java | 76 +++++++++++++++++-- .../group/TabletServerMetricGroup.java | 37 +++++++++ .../apache/fluss/server/kv/KvManagerTest.java | 48 ++++++++++++ .../kv/prewrite/KvPreWriteBufferTest.java | 30 ++++++++ .../observability/monitor-metrics.md | 10 +++ 8 files changed, 228 insertions(+), 6 deletions(-) diff --git a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java index 0219f3fe79f..b35d3edd805 100644 --- a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java +++ b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java @@ -258,6 +258,20 @@ 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"; + + /** + * Number of entries buffered in the pre-write buffers across all KV tablets in this server (Sum + * aggregation). + */ + public static final String KV_PRE_WRITE_BUFFER_ENTRY_COUNT = "kvPreWriteBufferEntryCount"; + // 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 = 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 14a14714bc5..fcec30ba7ac 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 @@ -231,6 +231,15 @@ private KvManager( this::getSharedBlockCachePinnedUsage, conf.get(ConfigOptions.KV_SHARED_BLOCK_CACHE_SIZE).getBytes()); } + tabletServerMetricGroup.setPreWriteBufferMetrics( + () -> + currentKvs.values().stream() + .mapToLong(KvTablet::kvPreWriteBufferMemoryUsageBytes) + .sum(), + () -> + currentKvs.values().stream() + .mapToInt(KvTablet::kvPreWriteBufferEntryCount) + .sum()); } private static RateLimiter createSharedRateLimiter(Configuration conf) { 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 61c5a2d5d7e..b085edf3bcc 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 @@ -853,6 +853,16 @@ FlushState getFlushState() { return inReadLock(kvLock, () -> flushState); } + /** Returns an instantaneous estimate of this kv tablet's pre-write buffer memory usage. */ + public long kvPreWriteBufferMemoryUsageBytes() { + return kvPreWriteBuffer.memoryUsageBytes(); + } + + /** Returns the number of entries held in this kv tablet's pre-write buffer. */ + public int kvPreWriteBufferEntryCount() { + return kvPreWriteBuffer.entryCount(); + } + @VisibleForTesting void setFlushState(FlushState state) { inWriteLock(kvLock, () -> flushState = state); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java index 18e922cf3f0..15c7ea27df5 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java @@ -89,6 +89,19 @@ @NotThreadSafe public class KvPreWriteBuffer { + /** + * Estimated JVM heap overhead of a single buffered entry besides its key/value payload bytes, + * covering the {@link KvEntry} object, its {@link Key} and {@link Value} wrappers, the byte + * array headers and the linked list node. + */ + private static final long PER_ENTRY_OVERHEAD_BYTES = 144; + + /** + * Estimated JVM heap overhead of a hash map node for an entry that is the latest version of its + * key in the buffer. + */ + private static final long PER_MAP_NODE_OVERHEAD_BYTES = 32; + // a mapping from the key to the kv-entry private final Map kvEntryMap = new HashMap<>(); @@ -105,6 +118,13 @@ public class KvPreWriteBuffer { // Accumulated byte size of entries not yet completed by a flush. private long pendingFlushBytes = 0; + // Estimated total memory footprint of the held entries, updated incrementally on the write + // path. + private volatile long memoryUsageBytes = 0; + + // Number of held entries. + private volatile int entryCount = 0; + public KvPreWriteBuffer(TabletServerMetricGroup serverMetricGroup) { truncateAsDuplicatedCount = serverMetricGroup.kvTruncateAsDuplicatedCount(); truncateAsErrorCount = serverMetricGroup.kvTruncateAsErrorCount(); @@ -181,8 +201,8 @@ private void doPut(ChangeType changeType, Key key, Value value, long lsn) { allKvEntries.addLast(kvEntry); // update the max lsn maxLogSequenceNumber = lsn; - // track accumulated bytes for flush budget gating - pendingFlushBytes += entryBytes(key, value); + // update the accounting for flush budget gating and metrics + addToAccounting(kvEntry); } /** @@ -226,8 +246,7 @@ public void truncateTo(long targetLogSequenceNumber, TruncateReason truncateReas + ", targetLogSequenceNumber=" + targetLogSequenceNumber); } - pendingFlushBytes -= entryBytes(entry.getKey(), entry.getValue()); - boolean removed = kvEntryMap.remove(entry.getKey(), entry); + boolean removed = removeFromMapAndAccounting(entry); // the removed entry is no longer the successor of its previous version; clear the // forward link so the truncated entry does not stay reachable through it if (entry.previousEntry != null) { @@ -251,6 +270,21 @@ public long pendingFlushBytes() { return pendingFlushBytes; } + /** Returns the number of entries currently held in this buffer. */ + public int entryCount() { + return entryCount; + } + + /** + * Returns an estimation of the total memory footprint of the entries currently held in this + * buffer, including the key/value payload bytes tracked by {@link #pendingFlushBytes()} and the + * per-entry JVM object overhead. This is an approximation for observability purposes, not an + * exact measurement. + */ + public long memoryUsageBytes() { + return memoryUsageBytes; + } + /** * Prepares a prefix of entries for asynchronous flush without removing them from the buffer. * @@ -286,8 +320,7 @@ public int completeFlush(PreparedFlush preparedFlush) { throw new IllegalStateException("Prepared flush entry is not in PREPARED state."); } entry.state = EntryState.FLUSHED; - pendingFlushBytes -= entryBytes(entry.getKey(), entry.getValue()); - kvEntryMap.remove(entry.getKey(), entry); + removeFromMapAndAccounting(entry); // the immediate successor is the only live referencer of a flushed entry; clearing // its reference makes the flushed entry (and, transitively, its older versions) // unreachable instead of being retained while no longer counted by pendingFlushBytes @@ -322,6 +355,37 @@ public void abortAllPrepared() { } } + /** + * Adds an entry to the incrementally maintained memory estimate and entry count. An entry + * without a previous version is the latest version of a new key and thus adds one map node. + */ + private void addToAccounting(KvEntry entry) { + long bytes = entryBytes(entry.getKey(), entry.getValue()); + pendingFlushBytes += bytes; + memoryUsageBytes += + bytes + + PER_ENTRY_OVERHEAD_BYTES + + (entry.previousEntry == null ? PER_MAP_NODE_OVERHEAD_BYTES : 0L); + entryCount++; + } + + /** + * Removes an entry from the key map and deducts it from the incrementally maintained memory + * estimate and entry count. Returns whether the entry was the latest version of its key and + * thus removed from the map. + */ + private boolean removeFromMapAndAccounting(KvEntry entry) { + long bytes = entryBytes(entry.getKey(), entry.getValue()); + pendingFlushBytes -= bytes; + boolean removedFromMap = kvEntryMap.remove(entry.getKey(), entry); + memoryUsageBytes -= + bytes + + PER_ENTRY_OVERHEAD_BYTES + + (removedFromMap ? PER_MAP_NODE_OVERHEAD_BYTES : 0L); + entryCount--; + return removedFromMap; + } + private static long entryBytes(Key key, Value value) { return (long) key.key.length + (value.value != null ? value.value.length : 0L); } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java index dac5562c10d..3d3c857ce28 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java @@ -90,6 +90,12 @@ public class TabletServerMetricGroup extends AbstractMetricGroup { private volatile long sharedWriteBufferCapacity; + /** Supplier for aggregated pre-write buffer memory usage, set by KvManager. */ + private volatile LongSupplier preWriteBufferMemoryUsageSupplier = () -> 0L; + + /** Supplier for aggregated pre-write buffer entry count, set by KvManager. */ + private volatile LongSupplier preWriteBufferEntryCountSupplier = () -> 0L; + public TabletServerMetricGroup( MetricRegistry registry, String clusterId, String rack, String hostname, int serverId) { super(registry, new String[] {clusterId, hostname, NAME}, null); @@ -151,6 +157,9 @@ public TabletServerMetricGroup( // Register server-level RocksDB aggregated metrics registerServerRocksDBMetrics(); + + // Register server-level pre-write buffer aggregated metrics + registerServerKvPreWriteBufferMetrics(); } /** @@ -215,6 +224,34 @@ public void setSharedWriteBufferMetrics(LongSupplier usageSupplier, long capacit this.sharedWriteBufferCapacity = capacity; } + /** + * Register server-level pre-write buffer aggregated metrics. These metrics aggregate the memory + * usage of the pre-write buffers of all KV tablets in this server. + */ + private void registerServerKvPreWriteBufferMetrics() { + gauge( + MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES, + () -> preWriteBufferMemoryUsageSupplier.getAsLong()); + gauge( + MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT, + () -> preWriteBufferEntryCountSupplier.getAsLong()); + } + + /** + * Sets the aggregated pre-write buffer metrics. Called by KvManager at construction time. + * + * @param memoryUsageSupplier supplier for the estimated memory usage of all pre-write buffers + * @param entryCountSupplier supplier for the total number of buffered entries across all + * pre-write buffers + */ + public void setPreWriteBufferMetrics( + LongSupplier memoryUsageSupplier, LongSupplier entryCountSupplier) { + this.preWriteBufferMemoryUsageSupplier = + checkNotNull(memoryUsageSupplier, "memoryUsageSupplier must not be null"); + this.preWriteBufferEntryCountSupplier = + checkNotNull(entryCountSupplier, "entryCountSupplier must not be null"); + } + /** * Registers gauges for the server-wide WAL memory pool used by primary key tables. Called once * by KvManager when creating the server buffer pool. diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java index e9c69e4fef6..63acf3578a1 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java @@ -282,6 +282,54 @@ void testSharedWriteBufferConfiguredThroughKvManagerCreateAndLoad() throws Excep .isEqualTo(capacity.getBytes()); } + @Test + void testPreWriteBufferServerLevelMetrics() throws Exception { + initTableBuckets(null); + assertThat( + gaugeValue( + TestingMetricGroups.TABLET_SERVER_METRICS, + MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) + .isEqualTo(0L); + assertThat( + gaugeValue( + TestingMetricGroups.TABLET_SERVER_METRICS, + MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) + .isEqualTo(0L); + + // write one kv record without flushing so that it stays in the pre-write buffer + KvTablet kvTablet = getOrCreateKv(tablePath1, null, tableBucket1); + KvRecordBatch kvRecordBatch = + kvRecordBatchFactory.ofRecords( + Collections.singletonList( + kvRecordFactory.ofRecord( + "key1".getBytes(), new Object[] {1, "a"}))); + kvTablet.putAsLeader(kvRecordBatch, null); + + assertThat( + gaugeValue( + TestingMetricGroups.TABLET_SERVER_METRICS, + MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) + .isEqualTo(1L); + assertThat( + gaugeValue( + TestingMetricGroups.TABLET_SERVER_METRICS, + MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) + .isPositive(); + + // the usage drops to zero once the buffered entries are flushed + flushAndWait(kvTablet, Long.MAX_VALUE); + assertThat( + gaugeValue( + TestingMetricGroups.TABLET_SERVER_METRICS, + MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) + .isEqualTo(0L); + assertThat( + gaugeValue( + TestingMetricGroups.TABLET_SERVER_METRICS, + MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) + .isEqualTo(0L); + } + @ParameterizedTest @MethodSource("partitionProvider") void testCreateKv(String partitionName) throws Exception { diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java index 45d983133a1..8aa09aabaec 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java @@ -468,6 +468,36 @@ void testPendingFlushBytesTracking() { assertThat(buffer.pendingFlushBytes()).isEqualTo(0); } + @Test + void testEstimatedMemoryUsage() { + KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + + assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); + assertThat(buffer.entryCount()).isEqualTo(0); + + // +key1(10 bytes), +key2(11 bytes), -key3(4 bytes): 25 payload bytes in total + bufferInsert(buffer, "key1", "value1", 1); + bufferInsert(buffer, "key2", "value22", 2); + bufferDelete(buffer, "key3", 3); + long payloadBytes = 25; + + // the estimation covers the payload bytes plus the per-entry object overhead + assertThat(buffer.memoryUsageBytes()).isGreaterThan(payloadBytes); + assertThat(buffer.entryCount()).isEqualTo(3); + + // flushing all entries releases the whole accounted usage + flushBuffer(buffer, Long.MAX_VALUE); + assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); + assertThat(buffer.entryCount()).isEqualTo(0); + + // truncating entries also releases their accounted usage + bufferInsert(buffer, "key1", "value1", 4); + assertThat(buffer.memoryUsageBytes()).isPositive(); + buffer.truncateTo(4, TruncateReason.ERROR); + assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); + assertThat(buffer.entryCount()).isEqualTo(0); + } + @Test void testCompleteFlushDetachesFlushedEntriesFromPreviousChain() { KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); diff --git a/website/docs/maintenance/observability/monitor-metrics.md b/website/docs/maintenance/observability/monitor-metrics.md index ee57382f523..d7dbe387fbf 100644 --- a/website/docs/maintenance/observability/monitor-metrics.md +++ b/website/docs/maintenance/observability/monitor-metrics.md @@ -588,6 +588,16 @@ Some metrics might not be exposed when using other JVM implementations (e.g. IBM preWriteBufferTruncateAsErrorPerSecond The number of kv pre-write buffer truncate due to the error happened when writing cdc to log per second. Meter + + + kvPreWriteBufferMemoryUsageBytes + Estimated total memory usage of the KV pre-write buffers across all KV tablets in this server (in bytes), including the key/value payload bytes and the per-entry object overhead. It is an approximation for observability, not an exact measurement. + Gauge + + + kvPreWriteBufferEntryCount + The number of entries buffered in the KV pre-write buffers across all KV tablets in this server. + Gauge kvWalMemoryPoolUsage From 0317392e7a173b4bc4e0c7e499203d7f1e096789 Mon Sep 17 00:00:00 2001 From: duankaixuan <1417048384@qq.com> Date: Wed, 16 Sep 2026 15:53:09 +0800 Subject: [PATCH 2/3] [server] Track pre-write buffer memory in a shared atomic ledger --- .../org/apache/fluss/server/kv/KvManager.java | 9 -- .../org/apache/fluss/server/kv/KvTablet.java | 13 +-- .../server/kv/prewrite/KvPreWriteBuffer.java | 45 ++++++++-- .../KvPreWriteBufferMemoryLedger.java | 83 +++++++++++++++++++ .../group/TabletServerMetricGroup.java | 29 ++----- .../apache/fluss/server/kv/KvManagerTest.java | 45 ++++------ .../kv/prewrite/KvPreWriteBufferTest.java | 39 ++++++++- 7 files changed, 184 insertions(+), 79 deletions(-) create mode 100644 fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferMemoryLedger.java 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 fcec30ba7ac..14a14714bc5 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 @@ -231,15 +231,6 @@ private KvManager( this::getSharedBlockCachePinnedUsage, conf.get(ConfigOptions.KV_SHARED_BLOCK_CACHE_SIZE).getBytes()); } - tabletServerMetricGroup.setPreWriteBufferMetrics( - () -> - currentKvs.values().stream() - .mapToLong(KvTablet::kvPreWriteBufferMemoryUsageBytes) - .sum(), - () -> - currentKvs.values().stream() - .mapToInt(KvTablet::kvPreWriteBufferEntryCount) - .sum()); } private static RateLimiter createSharedRateLimiter(Configuration conf) { 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 b085edf3bcc..a76a6398887 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 @@ -853,16 +853,6 @@ FlushState getFlushState() { return inReadLock(kvLock, () -> flushState); } - /** Returns an instantaneous estimate of this kv tablet's pre-write buffer memory usage. */ - public long kvPreWriteBufferMemoryUsageBytes() { - return kvPreWriteBuffer.memoryUsageBytes(); - } - - /** Returns the number of entries held in this kv tablet's pre-write buffer. */ - public int kvPreWriteBufferEntryCount() { - return kvPreWriteBuffer.entryCount(); - } - @VisibleForTesting void setFlushState(FlushState state) { inWriteLock(kvLock, () -> flushState = state); @@ -1356,6 +1346,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) { diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java index 15c7ea27df5..6a1fa4d7964 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java @@ -87,7 +87,7 @@ * head to tail, it will stop flush. */ @NotThreadSafe -public class KvPreWriteBuffer { +public class KvPreWriteBuffer implements AutoCloseable { /** * Estimated JVM heap overhead of a single buffered entry besides its key/value payload bytes, @@ -112,22 +112,29 @@ public class KvPreWriteBuffer { private final Counter truncateAsDuplicatedCount; private final Counter truncateAsErrorCount; + // The TabletServer-wide ledger shared by all pre-write buffers, updated atomically on the + // write path and serving as the single source of truth for metrics and backpressure. + private final KvPreWriteBufferMemoryLedger memoryLedger; + // the max LSN in the buffer private long maxLogSequenceNumber = -1; // Accumulated byte size of entries not yet completed by a flush. private long pendingFlushBytes = 0; - // Estimated total memory footprint of the held entries, updated incrementally on the write - // path. - private volatile long memoryUsageBytes = 0; + // Local accounting of this buffer, released to the shared ledger on close. Must be read + // under the kv write lock (or in single-threaded tests) to be exact. + private long memoryUsageBytes = 0; + + // Number of held entries, maintained together with the local memory accounting. + private int entryCount = 0; - // Number of held entries. - private volatile int entryCount = 0; + private boolean closed; public KvPreWriteBuffer(TabletServerMetricGroup serverMetricGroup) { truncateAsDuplicatedCount = serverMetricGroup.kvTruncateAsDuplicatedCount(); truncateAsErrorCount = serverMetricGroup.kvTruncateAsErrorCount(); + memoryLedger = serverMetricGroup.kvPreWriteBufferMemoryLedger(); } /** @@ -362,11 +369,13 @@ public void abortAllPrepared() { private void addToAccounting(KvEntry entry) { long bytes = entryBytes(entry.getKey(), entry.getValue()); pendingFlushBytes += bytes; - memoryUsageBytes += + long accountedBytes = bytes + PER_ENTRY_OVERHEAD_BYTES + (entry.previousEntry == null ? PER_MAP_NODE_OVERHEAD_BYTES : 0L); + memoryUsageBytes += accountedBytes; entryCount++; + memoryLedger.add(accountedBytes, 1); } /** @@ -378,14 +387,34 @@ private boolean removeFromMapAndAccounting(KvEntry entry) { long bytes = entryBytes(entry.getKey(), entry.getValue()); pendingFlushBytes -= bytes; boolean removedFromMap = kvEntryMap.remove(entry.getKey(), entry); - memoryUsageBytes -= + long accountedBytes = bytes + PER_ENTRY_OVERHEAD_BYTES + (removedFromMap ? PER_MAP_NODE_OVERHEAD_BYTES : 0L); + memoryUsageBytes -= accountedBytes; entryCount--; + memoryLedger.subtract(accountedBytes, 1); return removedFromMap; } + /** + * Closes the buffer and releases its remaining accounting to the shared ledger. Must be called + * under the kv write lock so the local accounting values are exact. Idempotent. + */ + @Override + public void close() { + if (closed) { + return; + } + closed = true; + memoryLedger.subtract(memoryUsageBytes, entryCount); + memoryUsageBytes = 0; + entryCount = 0; + allKvEntries.clear(); + kvEntryMap.clear(); + maxLogSequenceNumber = -1; + } + private static long entryBytes(Key key, Value value) { return (long) key.key.length + (value.value != null ? value.value.length : 0L); } diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferMemoryLedger.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferMemoryLedger.java new file mode 100644 index 00000000000..48536969abc --- /dev/null +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferMemoryLedger.java @@ -0,0 +1,83 @@ +/* + * 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.prewrite; + +import org.apache.fluss.annotation.Internal; + +import javax.annotation.concurrent.ThreadSafe; + +import java.util.concurrent.atomic.AtomicLong; + +import static org.apache.fluss.utils.Preconditions.checkArgument; + +/** + * TabletServer-wide memory ledger shared by all KV pre-write buffers. It atomically tracks the + * total estimated memory usage and the total number of entries held across all buffers, and serves + * as the single source of truth for both the exposed metrics and future backpressure decisions. + * + *

Each buffer reports its accounting deltas to this ledger on the write path through {@link + * #add} and {@link #subtract}, so a metric read is a single atomic snapshot instead of a sum of + * non-atomic samples collected from multiple buffers. + * + *

The memory usage covers the key/value payload bytes plus a per-entry object overhead + * approximation, reflecting the real retained heap of the buffered entries. + */ +@Internal +@ThreadSafe +public final class KvPreWriteBufferMemoryLedger { + + private final AtomicLong memoryUsageBytes = new AtomicLong(); + + private final AtomicLong entryCount = new AtomicLong(); + + /** + * Adds the given amount of memory usage and entries to the ledger. Called when entries are + * appended to a pre-write buffer. + */ + public void add(long memoryBytes, int entries) { + checkArgument(memoryBytes >= 0, "The added memory bytes must not be negative."); + checkArgument(entries >= 0, "The added entry count must not be negative."); + memoryUsageBytes.addAndGet(memoryBytes); + entryCount.addAndGet(entries); + } + + /** + * Subtracts the given amount of memory usage and entries from the ledger. Called when entries + * leave a pre-write buffer by flushing or truncation, or when a buffer is closed. + */ + public void subtract(long memoryBytes, int entries) { + checkArgument(memoryBytes >= 0, "The subtracted memory bytes must not be negative."); + checkArgument(entries >= 0, "The subtracted entry count must not be negative."); + memoryUsageBytes.addAndGet(-memoryBytes); + entryCount.addAndGet(-entries); + } + + /** + * Returns the total estimated memory usage across all pre-write buffers in bytes, including the + * key/value payload bytes and the per-entry object overhead. This is an approximation for + * observability purposes, not an exact measurement. + */ + public long memoryUsageBytes() { + return memoryUsageBytes.get(); + } + + /** Returns the total number of entries held across all pre-write buffers. */ + public long entryCount() { + return entryCount.get(); + } +} diff --git a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java index 3d3c857ce28..c62e864c3bb 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java @@ -30,6 +30,7 @@ import org.apache.fluss.metrics.ThreadSafeSimpleCounter; import org.apache.fluss.metrics.groups.AbstractMetricGroup; import org.apache.fluss.metrics.registry.MetricRegistry; +import org.apache.fluss.server.kv.prewrite.KvPreWriteBufferMemoryLedger; import org.apache.fluss.server.kv.rocksdb.RocksDBStatistics; import java.util.Map; @@ -90,11 +91,9 @@ public class TabletServerMetricGroup extends AbstractMetricGroup { private volatile long sharedWriteBufferCapacity; - /** Supplier for aggregated pre-write buffer memory usage, set by KvManager. */ - private volatile LongSupplier preWriteBufferMemoryUsageSupplier = () -> 0L; - - /** Supplier for aggregated pre-write buffer entry count, set by KvManager. */ - private volatile LongSupplier preWriteBufferEntryCountSupplier = () -> 0L; + /** Ledger shared by all KV pre-write buffers, serving as the single accounting source. */ + private final KvPreWriteBufferMemoryLedger kvPreWriteBufferMemoryLedger = + new KvPreWriteBufferMemoryLedger(); public TabletServerMetricGroup( MetricRegistry registry, String clusterId, String rack, String hostname, int serverId) { @@ -231,25 +230,15 @@ public void setSharedWriteBufferMetrics(LongSupplier usageSupplier, long capacit private void registerServerKvPreWriteBufferMetrics() { gauge( MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES, - () -> preWriteBufferMemoryUsageSupplier.getAsLong()); + kvPreWriteBufferMemoryLedger::memoryUsageBytes); gauge( MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT, - () -> preWriteBufferEntryCountSupplier.getAsLong()); + kvPreWriteBufferMemoryLedger::entryCount); } - /** - * Sets the aggregated pre-write buffer metrics. Called by KvManager at construction time. - * - * @param memoryUsageSupplier supplier for the estimated memory usage of all pre-write buffers - * @param entryCountSupplier supplier for the total number of buffered entries across all - * pre-write buffers - */ - public void setPreWriteBufferMetrics( - LongSupplier memoryUsageSupplier, LongSupplier entryCountSupplier) { - this.preWriteBufferMemoryUsageSupplier = - checkNotNull(memoryUsageSupplier, "memoryUsageSupplier must not be null"); - this.preWriteBufferEntryCountSupplier = - checkNotNull(entryCountSupplier, "entryCountSupplier must not be null"); + /** Returns the memory ledger shared by all KV pre-write buffers of this server. */ + public KvPreWriteBufferMemoryLedger kvPreWriteBufferMemoryLedger() { + return kvPreWriteBufferMemoryLedger; } /** diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java index 63acf3578a1..4333478dc7d 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java @@ -285,16 +285,11 @@ void testSharedWriteBufferConfiguredThroughKvManagerCreateAndLoad() throws Excep @Test void testPreWriteBufferServerLevelMetrics() throws Exception { initTableBuckets(null); - assertThat( - gaugeValue( - TestingMetricGroups.TABLET_SERVER_METRICS, - MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) - .isEqualTo(0L); - assertThat( - gaugeValue( - TestingMetricGroups.TABLET_SERVER_METRICS, - MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) - .isEqualTo(0L); + TabletServerMetricGroup metricGroup = TestingMetricGroups.TABLET_SERVER_METRICS; + long memoryUsageBefore = + gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES); + long entryCountBefore = + gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT); // write one kv record without flushing so that it stays in the pre-write buffer KvTablet kvTablet = getOrCreateKv(tablePath1, null, tableBucket1); @@ -305,29 +300,17 @@ void testPreWriteBufferServerLevelMetrics() throws Exception { "key1".getBytes(), new Object[] {1, "a"}))); kvTablet.putAsLeader(kvRecordBatch, null); - assertThat( - gaugeValue( - TestingMetricGroups.TABLET_SERVER_METRICS, - MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) - .isEqualTo(1L); - assertThat( - gaugeValue( - TestingMetricGroups.TABLET_SERVER_METRICS, - MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) - .isPositive(); + assertThat(gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) + .isEqualTo(entryCountBefore + 1); + assertThat(gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) + .isGreaterThan(memoryUsageBefore); - // the usage drops to zero once the buffered entries are flushed + // the accounting returns to its previous values once the buffered entry is flushed flushAndWait(kvTablet, Long.MAX_VALUE); - assertThat( - gaugeValue( - TestingMetricGroups.TABLET_SERVER_METRICS, - MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) - .isEqualTo(0L); - assertThat( - gaugeValue( - TestingMetricGroups.TABLET_SERVER_METRICS, - MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) - .isEqualTo(0L); + assertThat(gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) + .isEqualTo(entryCountBefore); + assertThat(gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) + .isEqualTo(memoryUsageBefore); } @ParameterizedTest diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java index 8aa09aabaec..a277f02e11b 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java @@ -17,8 +17,10 @@ package org.apache.fluss.server.kv.prewrite; +import org.apache.fluss.metrics.registry.NOPMetricRegistry; import org.apache.fluss.server.kv.prewrite.KvPreWriteBuffer.PreparedFlush; import org.apache.fluss.server.kv.prewrite.KvPreWriteBuffer.TruncateReason; +import org.apache.fluss.server.metrics.group.TabletServerMetricGroup; import org.apache.fluss.server.metrics.group.TestingMetricGroups; import org.junit.jupiter.api.Test; @@ -470,10 +472,14 @@ void testPendingFlushBytesTracking() { @Test void testEstimatedMemoryUsage() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + TabletServerMetricGroup metricGroup = + new TabletServerMetricGroup(NOPMetricRegistry.INSTANCE, "fluss", "rack", "host", 0); + KvPreWriteBuffer buffer = new KvPreWriteBuffer(metricGroup); assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); assertThat(buffer.entryCount()).isEqualTo(0); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); // +key1(10 bytes), +key2(11 bytes), -key3(4 bytes): 25 payload bytes in total bufferInsert(buffer, "key1", "value1", 1); @@ -484,11 +490,17 @@ void testEstimatedMemoryUsage() { // the estimation covers the payload bytes plus the per-entry object overhead assertThat(buffer.memoryUsageBytes()).isGreaterThan(payloadBytes); assertThat(buffer.entryCount()).isEqualTo(3); + // the shared ledger mirrors the local accounting of the buffer + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()) + .isEqualTo(buffer.memoryUsageBytes()); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(3L); // flushing all entries releases the whole accounted usage flushBuffer(buffer, Long.MAX_VALUE); assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); assertThat(buffer.entryCount()).isEqualTo(0); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); // truncating entries also releases their accounted usage bufferInsert(buffer, "key1", "value1", 4); @@ -496,6 +508,31 @@ void testEstimatedMemoryUsage() { buffer.truncateTo(4, TruncateReason.ERROR); assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); assertThat(buffer.entryCount()).isEqualTo(0); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); + } + + @Test + void testCloseReleasesAccounting() { + TabletServerMetricGroup metricGroup = + new TabletServerMetricGroup(NOPMetricRegistry.INSTANCE, "fluss", "rack", "host", 0); + KvPreWriteBuffer buffer = new KvPreWriteBuffer(metricGroup); + bufferInsert(buffer, "key1", "value1", 0); + bufferInsert(buffer, "key2", "value2", 1); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isPositive(); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(2L); + + // closing releases the remaining accounting to the shared ledger exactly once + buffer.close(); + assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); + assertThat(buffer.entryCount()).isEqualTo(0); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); + + // closing again is idempotent and must not over-release the ledger + buffer.close(); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); + assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); } @Test From 323b774fc5ee1e876a6dd26f5f90a354e56e343a Mon Sep 17 00:00:00 2001 From: duankaixuan <1417048384@qq.com> Date: Thu, 17 Sep 2026 16:14:37 +0800 Subject: [PATCH 3/3] [server] Address review: shared counter ownership and map-node accounting fix --- .../org/apache/fluss/metrics/MetricNames.java | 6 - .../org/apache/fluss/server/kv/KvManager.java | 8 + .../org/apache/fluss/server/kv/KvTablet.java | 13 +- .../server/kv/prewrite/KvPreWriteBuffer.java | 194 ++++++++++-------- .../KvPreWriteBufferMemoryLedger.java | 83 -------- .../group/TabletServerMetricGroup.java | 33 ++- .../apache/fluss/server/kv/KvManagerTest.java | 58 ++++-- .../apache/fluss/server/kv/KvTabletTest.java | 8 + .../kv/prewrite/KvPreWriteBufferTest.java | 157 ++++++++++---- .../observability/monitor-metrics.md | 9 +- 10 files changed, 312 insertions(+), 257 deletions(-) delete mode 100644 fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferMemoryLedger.java diff --git a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java index b35d3edd805..508b243d800 100644 --- a/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java +++ b/fluss-common/src/main/java/org/apache/fluss/metrics/MetricNames.java @@ -266,12 +266,6 @@ public class MetricNames { public static final String KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES = "kvPreWriteBufferMemoryUsageBytes"; - /** - * Number of entries buffered in the pre-write buffers across all KV tablets in this server (Sum - * aggregation). - */ - public static final String KV_PRE_WRITE_BUFFER_ENTRY_COUNT = "kvPreWriteBufferEntryCount"; - // 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 = 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 14a14714bc5..b1db6ed343a 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 @@ -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; @@ -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; @@ -225,6 +229,8 @@ private KvManager( throw e; } this.kvFlushScheduler = createdFlushScheduler; + tabletServerMetricGroup.setKvPreWriteBufferMemoryUsageMetrics( + kvPreWriteBufferMemoryUsageBytes::get); if (sharedBlockCache != null) { tabletServerMetricGroup.setSharedBlockCacheMetrics( this::getSharedBlockCacheUsage, @@ -459,6 +465,7 @@ public KvTablet getOrCreateKv( sharedRocksDBRateLimiter, sharedBlockCache, sharedWriteBufferManager, + kvPreWriteBufferMemoryUsageBytes, kvFlushScheduler, flushCompleteListener, autoIncrementManager, @@ -582,6 +589,7 @@ public KvTablet loadKv( sharedRocksDBRateLimiter, sharedBlockCache, sharedWriteBufferManager, + kvPreWriteBufferMemoryUsageBytes, kvFlushScheduler, flushCompleteListener, autoIncrementManager, 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 a76a6398887..4fd15ff0569 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 @@ -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; @@ -203,6 +204,7 @@ private KvTablet( ValueEncoder valueEncoder, ValueDecoder valueDecoder, @Nullable RocksDBStatistics rocksDBStatistics, + AtomicLong sharedPreWriteBufferMemoryUsageBytes, KvFlushScheduler kvFlushScheduler, boolean closeFlushScheduler, @Nullable Runnable flushCompleteListener, @@ -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; @@ -299,6 +302,7 @@ public static KvTablet create( sharedRateLimiter, null, null, + new AtomicLong(), new KvFlushScheduler(serverConf), true, null, @@ -324,6 +328,7 @@ static KvTablet create( RateLimiter sharedRateLimiter, @Nullable Cache sharedBlockCache, @Nullable WriteBufferManager sharedWriteBufferManager, + AtomicLong sharedPreWriteBufferMemoryUsageBytes, KvFlushScheduler kvFlushScheduler, @Nullable Runnable flushCompleteListener, AutoIncrementManager autoIncrementManager, @@ -347,6 +352,7 @@ static KvTablet create( sharedRateLimiter, sharedBlockCache, sharedWriteBufferManager, + sharedPreWriteBufferMemoryUsageBytes, kvFlushScheduler, false, flushCompleteListener, @@ -371,6 +377,7 @@ public static KvTablet create( ChangelogImage changelogImage, RateLimiter sharedRateLimiter, @Nullable Cache sharedBlockCache, + AtomicLong sharedPreWriteBufferMemoryUsageBytes, KvFlushScheduler kvFlushScheduler, @Nullable Runnable flushCompleteListener, AutoIncrementManager autoIncrementManager, @@ -394,6 +401,7 @@ public static KvTablet create( sharedRateLimiter, sharedBlockCache, null, + sharedPreWriteBufferMemoryUsageBytes, kvFlushScheduler, false, flushCompleteListener, @@ -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, @@ -487,6 +496,7 @@ private static KvTablet create( valueEncoder, valueDecoder, rocksDBStatistics, + sharedPreWriteBufferMemoryUsageBytes, kvFlushScheduler, closeFlushScheduler, flushCompleteListener, @@ -532,6 +542,7 @@ public static KvTablet create( sharedRateLimiter, null, null, + new AtomicLong(), new KvFlushScheduler(serverConf), true, null, diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java index 6a1fa4d7964..9de3c56b842 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java @@ -37,6 +37,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.concurrent.atomic.AtomicLong; import static org.apache.fluss.utils.Preconditions.checkArgument; import static org.apache.fluss.utils.UnsafeUtils.BYTE_ARRAY_BASE_OFFSET; @@ -112,9 +113,9 @@ public class KvPreWriteBuffer implements AutoCloseable { private final Counter truncateAsDuplicatedCount; private final Counter truncateAsErrorCount; - // The TabletServer-wide ledger shared by all pre-write buffers, updated atomically on the - // write path and serving as the single source of truth for metrics and backpressure. - private final KvPreWriteBufferMemoryLedger memoryLedger; + // The server-wide counter shared by all pre-write buffers, updated atomically on the write + // path and serving as the single source of truth for the memory usage metric. + private final AtomicLong memoryUsageBytesCounter; // the max LSN in the buffer private long maxLogSequenceNumber = -1; @@ -122,19 +123,17 @@ public class KvPreWriteBuffer implements AutoCloseable { // Accumulated byte size of entries not yet completed by a flush. private long pendingFlushBytes = 0; - // Local accounting of this buffer, released to the shared ledger on close. Must be read + // Local accounting of this buffer, released to the shared counter on close. Must be read // under the kv write lock (or in single-threaded tests) to be exact. private long memoryUsageBytes = 0; - // Number of held entries, maintained together with the local memory accounting. - private int entryCount = 0; - private boolean closed; - public KvPreWriteBuffer(TabletServerMetricGroup serverMetricGroup) { + public KvPreWriteBuffer( + TabletServerMetricGroup serverMetricGroup, AtomicLong memoryUsageBytesCounter) { truncateAsDuplicatedCount = serverMetricGroup.kvTruncateAsDuplicatedCount(); truncateAsErrorCount = serverMetricGroup.kvTruncateAsErrorCount(); - memoryLedger = serverMetricGroup.kvPreWriteBufferMemoryLedger(); + this.memoryUsageBytesCounter = memoryUsageBytesCounter; } /** @@ -228,6 +227,10 @@ private void doPut(ChangeType changeType, Key key, Value value, long lsn) { * Truncate the buffer to the given log sequence number so that it only contains key-value pairs * whose log sequence number is less than the given log sequence number. * + *

The bytes released by the truncated entries are applied to the shared counter in one + * update, so a truncation that fails partway (e.g. on a prepared entry) still releases the + * accounting of the entries removed before the failure. + * * @param targetLogSequenceNumber the lower bound of the log sequence number truncated to. * @param truncateReason the reason to truncate */ @@ -238,38 +241,54 @@ public void truncateTo(long targetLogSequenceNumber, TruncateReason truncateReas truncateAsErrorCount.inc(); } - Iterator descIter = allKvEntries.descendingIterator(); - while (descIter.hasNext()) { - KvEntry entry = descIter.next(); - if (entry.getLogSequenceNumber() < targetLogSequenceNumber) { - maxLogSequenceNumber = entry.logSequenceNumber; - break; - } - descIter.remove(); - if (entry.state == EntryState.PREPARED) { - throw new IllegalStateException( - "Cannot truncate prepared pre-write entry. logSequenceNumber=" - + entry.getLogSequenceNumber() - + ", targetLogSequenceNumber=" - + targetLogSequenceNumber); + // Net bytes released by the entries this call actually removes; applied to the shared + // counter in one update below. + long netReleasedBytes = 0; + try { + Iterator descIter = allKvEntries.descendingIterator(); + while (descIter.hasNext()) { + KvEntry entry = descIter.next(); + if (entry.getLogSequenceNumber() < targetLogSequenceNumber) { + maxLogSequenceNumber = entry.logSequenceNumber; + break; + } + descIter.remove(); + if (entry.state == EntryState.PREPARED) { + throw new IllegalStateException( + "Cannot truncate prepared pre-write entry. logSequenceNumber=" + + entry.getLogSequenceNumber() + + ", targetLogSequenceNumber=" + + targetLogSequenceNumber); + } + boolean removedFromMap = removeFromMapAndAccounting(entry); + netReleasedBytes += entryAccountedBytes(entry, removedFromMap); + // the removed entry is no longer the successor of its previous version; clear + // the forward link so the truncated entry does not stay reachable through it + if (entry.previousEntry != null) { + entry.previousEntry.nextEntry = null; + } + // if the latest entry is removed, we need to rollback the previous entry to + // the map + if (removedFromMap) { + KvEntry previousEntry = previousEntryInBuffer(entry.previousEntry); + if (previousEntry != null) { + kvEntryMap.put(entry.getKey(), previousEntry); + // reinstating the older version re-occupies the key's map node, so + // restore its accounting as a negative release in the same batched + // update + memoryUsageBytes += PER_MAP_NODE_OVERHEAD_BYTES; + netReleasedBytes -= PER_MAP_NODE_OVERHEAD_BYTES; + } + } } - boolean removed = removeFromMapAndAccounting(entry); - // the removed entry is no longer the successor of its previous version; clear the - // forward link so the truncated entry does not stay reachable through it - if (entry.previousEntry != null) { - entry.previousEntry.nextEntry = null; + if (!descIter.hasNext()) { + maxLogSequenceNumber = -1; } - // if the latest entry is removed, we need to rollback the previous entry to the map - if (removed) { - KvEntry previousEntry = previousEntryInBuffer(entry.previousEntry); - if (previousEntry != null) { - kvEntryMap.put(entry.getKey(), previousEntry); - } + } finally { + if (netReleasedBytes != 0) { + memoryUsageBytesCounter.addAndGet(-netReleasedBytes); } } - if (!descIter.hasNext()) { - maxLogSequenceNumber = -1; - } } /** Returns the accumulated byte size of all entries waiting to be flushed. */ @@ -277,11 +296,6 @@ public long pendingFlushBytes() { return pendingFlushBytes; } - /** Returns the number of entries currently held in this buffer. */ - public int entryCount() { - return entryCount; - } - /** * Returns an estimation of the total memory footprint of the entries currently held in this * buffer, including the key/value payload bytes tracked by {@link #pendingFlushBytes()} and the @@ -316,23 +330,38 @@ public PreparedFlush prepareFlush(long exclusiveUpToLogSequenceNumber) { return new PreparedFlush(exclusiveUpToLogSequenceNumber, entries, rowCountDiff); } - /** Completes a prepared async flush and removes flushed entries from the buffer. */ + /** + * Completes a prepared async flush and removes flushed entries from the buffer. The bytes + * released by the removed entries are applied to the shared counter in one update, so a partial + * failure still releases the accounting of the entries removed before the failure. + */ public int completeFlush(PreparedFlush preparedFlush) { - for (KvEntry entry : preparedFlush.entries) { - KvEntry first = allKvEntries.removeFirst(); - if (first != entry) { - throw new IllegalStateException("Prepared flush entries are no longer a prefix."); - } - if (entry.state != EntryState.PREPARED) { - throw new IllegalStateException("Prepared flush entry is not in PREPARED state."); + long releasedBytes = 0; + try { + for (KvEntry entry : preparedFlush.entries) { + KvEntry first = allKvEntries.removeFirst(); + if (first != entry) { + throw new IllegalStateException( + "Prepared flush entries are no longer a prefix."); + } + if (entry.state != EntryState.PREPARED) { + throw new IllegalStateException( + "Prepared flush entry is not in PREPARED state."); + } + entry.state = EntryState.FLUSHED; + boolean removedFromMap = removeFromMapAndAccounting(entry); + releasedBytes += entryAccountedBytes(entry, removedFromMap); + // the immediate successor is the only live referencer of a flushed entry; + // clearing its reference makes the flushed entry (and, transitively, its older + // versions) unreachable instead of being retained while no longer counted by + // pendingFlushBytes + if (entry.nextEntry != null) { + entry.nextEntry.previousEntry = null; + } } - entry.state = EntryState.FLUSHED; - removeFromMapAndAccounting(entry); - // the immediate successor is the only live referencer of a flushed entry; clearing - // its reference makes the flushed entry (and, transitively, its older versions) - // unreachable instead of being retained while no longer counted by pendingFlushBytes - if (entry.nextEntry != null) { - entry.nextEntry.previousEntry = null; + } finally { + if (releasedBytes != 0) { + memoryUsageBytesCounter.addAndGet(-releasedBytes); } } if (allKvEntries.isEmpty()) { @@ -363,43 +392,43 @@ public void abortAllPrepared() { } /** - * Adds an entry to the incrementally maintained memory estimate and entry count. An entry - * without a previous version is the latest version of a new key and thus adds one map node. + * Adds an entry to the incrementally maintained memory accounting. An entry without a previous + * version is the latest version of a new key and thus adds one map node. The put path appends + * entries one by one, so each accounted entry is published to the shared counter directly. */ private void addToAccounting(KvEntry entry) { - long bytes = entryBytes(entry.getKey(), entry.getValue()); - pendingFlushBytes += bytes; - long accountedBytes = - bytes - + PER_ENTRY_OVERHEAD_BYTES - + (entry.previousEntry == null ? PER_MAP_NODE_OVERHEAD_BYTES : 0L); + pendingFlushBytes += entryBytes(entry.getKey(), entry.getValue()); + long accountedBytes = entryAccountedBytes(entry, entry.previousEntry == null); memoryUsageBytes += accountedBytes; - entryCount++; - memoryLedger.add(accountedBytes, 1); + memoryUsageBytesCounter.addAndGet(accountedBytes); } /** - * Removes an entry from the key map and deducts it from the incrementally maintained memory - * estimate and entry count. Returns whether the entry was the latest version of its key and - * thus removed from the map. + * Removes an entry from the key map and deducts it from the local accounting only. The shared + * counter is not updated here: callers accumulate the bytes released by the whole flush or + * truncate operation and apply them to the shared counter in one update. Returns whether the + * entry was the latest version of its key and thus removed from the map. */ private boolean removeFromMapAndAccounting(KvEntry entry) { - long bytes = entryBytes(entry.getKey(), entry.getValue()); - pendingFlushBytes -= bytes; + pendingFlushBytes -= entryBytes(entry.getKey(), entry.getValue()); boolean removedFromMap = kvEntryMap.remove(entry.getKey(), entry); - long accountedBytes = - bytes - + PER_ENTRY_OVERHEAD_BYTES - + (removedFromMap ? PER_MAP_NODE_OVERHEAD_BYTES : 0L); - memoryUsageBytes -= accountedBytes; - entryCount--; - memoryLedger.subtract(accountedBytes, 1); + memoryUsageBytes -= entryAccountedBytes(entry, removedFromMap); return removedFromMap; } /** - * Closes the buffer and releases its remaining accounting to the shared ledger. Must be called - * under the kv write lock so the local accounting values are exact. Idempotent. + * Returns the accounting of one entry: its key/value payload bytes, the per-entry object + * overhead, and the map-node overhead if the entry holds the key's map node. + */ + private static long entryAccountedBytes(KvEntry entry, boolean holdsMapNode) { + return entryBytes(entry.getKey(), entry.getValue()) + + PER_ENTRY_OVERHEAD_BYTES + + (holdsMapNode ? PER_MAP_NODE_OVERHEAD_BYTES : 0L); + } + + /** + * Closes the buffer and releases its remaining accounting to the shared counter. Must be called + * under the kv write lock so the local accounting value is exact. Idempotent. */ @Override public void close() { @@ -407,9 +436,8 @@ public void close() { return; } closed = true; - memoryLedger.subtract(memoryUsageBytes, entryCount); + memoryUsageBytesCounter.addAndGet(-memoryUsageBytes); memoryUsageBytes = 0; - entryCount = 0; allKvEntries.clear(); kvEntryMap.clear(); maxLogSequenceNumber = -1; diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferMemoryLedger.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferMemoryLedger.java deleted file mode 100644 index 48536969abc..00000000000 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferMemoryLedger.java +++ /dev/null @@ -1,83 +0,0 @@ -/* - * 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.prewrite; - -import org.apache.fluss.annotation.Internal; - -import javax.annotation.concurrent.ThreadSafe; - -import java.util.concurrent.atomic.AtomicLong; - -import static org.apache.fluss.utils.Preconditions.checkArgument; - -/** - * TabletServer-wide memory ledger shared by all KV pre-write buffers. It atomically tracks the - * total estimated memory usage and the total number of entries held across all buffers, and serves - * as the single source of truth for both the exposed metrics and future backpressure decisions. - * - *

Each buffer reports its accounting deltas to this ledger on the write path through {@link - * #add} and {@link #subtract}, so a metric read is a single atomic snapshot instead of a sum of - * non-atomic samples collected from multiple buffers. - * - *

The memory usage covers the key/value payload bytes plus a per-entry object overhead - * approximation, reflecting the real retained heap of the buffered entries. - */ -@Internal -@ThreadSafe -public final class KvPreWriteBufferMemoryLedger { - - private final AtomicLong memoryUsageBytes = new AtomicLong(); - - private final AtomicLong entryCount = new AtomicLong(); - - /** - * Adds the given amount of memory usage and entries to the ledger. Called when entries are - * appended to a pre-write buffer. - */ - public void add(long memoryBytes, int entries) { - checkArgument(memoryBytes >= 0, "The added memory bytes must not be negative."); - checkArgument(entries >= 0, "The added entry count must not be negative."); - memoryUsageBytes.addAndGet(memoryBytes); - entryCount.addAndGet(entries); - } - - /** - * Subtracts the given amount of memory usage and entries from the ledger. Called when entries - * leave a pre-write buffer by flushing or truncation, or when a buffer is closed. - */ - public void subtract(long memoryBytes, int entries) { - checkArgument(memoryBytes >= 0, "The subtracted memory bytes must not be negative."); - checkArgument(entries >= 0, "The subtracted entry count must not be negative."); - memoryUsageBytes.addAndGet(-memoryBytes); - entryCount.addAndGet(-entries); - } - - /** - * Returns the total estimated memory usage across all pre-write buffers in bytes, including the - * key/value payload bytes and the per-entry object overhead. This is an approximation for - * observability purposes, not an exact measurement. - */ - public long memoryUsageBytes() { - return memoryUsageBytes.get(); - } - - /** Returns the total number of entries held across all pre-write buffers. */ - public long entryCount() { - return entryCount.get(); - } -} diff --git a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java index c62e864c3bb..f3de1b34c03 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java @@ -30,7 +30,6 @@ import org.apache.fluss.metrics.ThreadSafeSimpleCounter; import org.apache.fluss.metrics.groups.AbstractMetricGroup; import org.apache.fluss.metrics.registry.MetricRegistry; -import org.apache.fluss.server.kv.prewrite.KvPreWriteBufferMemoryLedger; import org.apache.fluss.server.kv.rocksdb.RocksDBStatistics; import java.util.Map; @@ -91,9 +90,8 @@ public class TabletServerMetricGroup extends AbstractMetricGroup { private volatile long sharedWriteBufferCapacity; - /** Ledger shared by all KV pre-write buffers, serving as the single accounting source. */ - private final KvPreWriteBufferMemoryLedger kvPreWriteBufferMemoryLedger = - new KvPreWriteBufferMemoryLedger(); + /** Supplier for the server-wide pre-write buffer memory usage, set by KvManager. */ + private volatile LongSupplier kvPreWriteBufferMemoryUsageSupplier = () -> 0L; public TabletServerMetricGroup( MetricRegistry registry, String clusterId, String rack, String hostname, int serverId) { @@ -157,8 +155,9 @@ public TabletServerMetricGroup( // Register server-level RocksDB aggregated metrics registerServerRocksDBMetrics(); - // Register server-level pre-write buffer aggregated metrics - registerServerKvPreWriteBufferMetrics(); + gauge( + MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES, + () -> kvPreWriteBufferMemoryUsageSupplier.getAsLong()); } /** @@ -224,21 +223,15 @@ public void setSharedWriteBufferMetrics(LongSupplier usageSupplier, long capacit } /** - * Register server-level pre-write buffer aggregated metrics. These metrics aggregate the memory - * usage of the pre-write buffers of all KV tablets in this server. + * Sets the supplier for the server-wide pre-write buffer memory usage gauge. Called by + * KvManager, which owns the shared accounting counter. + * + * @param usageSupplier supplier for the current total estimated memory usage of all KV + * pre-write buffers in bytes */ - private void registerServerKvPreWriteBufferMetrics() { - gauge( - MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES, - kvPreWriteBufferMemoryLedger::memoryUsageBytes); - gauge( - MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT, - kvPreWriteBufferMemoryLedger::entryCount); - } - - /** Returns the memory ledger shared by all KV pre-write buffers of this server. */ - public KvPreWriteBufferMemoryLedger kvPreWriteBufferMemoryLedger() { - return kvPreWriteBufferMemoryLedger; + public void setKvPreWriteBufferMemoryUsageMetrics(LongSupplier usageSupplier) { + this.kvPreWriteBufferMemoryUsageSupplier = + checkNotNull(usageSupplier, "usageSupplier must not be null"); } /** diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java index 4333478dc7d..da87031f3e2 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/KvManagerTest.java @@ -75,10 +75,12 @@ import java.util.Collections; import java.util.List; import java.util.Optional; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; import static org.apache.fluss.compression.ArrowCompressionInfo.DEFAULT_COMPRESSION; import static org.apache.fluss.record.TestData.DATA1_SCHEMA_PK; @@ -288,27 +290,49 @@ void testPreWriteBufferServerLevelMetrics() throws Exception { TabletServerMetricGroup metricGroup = TestingMetricGroups.TABLET_SERVER_METRICS; long memoryUsageBefore = gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES); - long entryCountBefore = - gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT); - // write one kv record without flushing so that it stays in the pre-write buffer KvTablet kvTablet = getOrCreateKv(tablePath1, null, tableBucket1); - KvRecordBatch kvRecordBatch = - kvRecordBatchFactory.ofRecords( - Collections.singletonList( - kvRecordFactory.ofRecord( - "key1".getBytes(), new Object[] {1, "a"}))); - kvTablet.putAsLeader(kvRecordBatch, null); - - assertThat(gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) - .isEqualTo(entryCountBefore + 1); - assertThat(gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) - .isGreaterThan(memoryUsageBefore); + // block the async flush right before its native write, so the buffered entry stays + // accounted while we assert on the gauge + CountDownLatch flushEnteredNativeWrite = new CountDownLatch(1); + CountDownLatch releaseNativeWrite = new CountDownLatch(1); + kvTablet.setBeforeNativeWrite( + () -> { + flushEnteredNativeWrite.countDown(); + try { + releaseNativeWrite.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }); - // the accounting returns to its previous values once the buffered entry is flushed + try { + // write one kv record; it stays in the pre-write buffer until the flush completes + KvRecordBatch kvRecordBatch = + kvRecordBatchFactory.ofRecords( + Collections.singletonList( + kvRecordFactory.ofRecord( + "key1".getBytes(), new Object[] {1, "a"}))); + kvTablet.putAsLeader(kvRecordBatch, null); + + // the put path itself does not schedule a flush; request one explicitly so the + // async flush reaches its native write while the entry is still accounted + AtomicReference flushFailure = new AtomicReference<>(); + kvTablet.requestFlush(kvTablet.localLogEndOffset(), flushFailure::set); + + // wait until the async flush reaches its native write; the entry is still accounted + flushEnteredNativeWrite.await(); + assertThat(flushFailure.get()).isNull(); + assertThat(gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) + .isGreaterThan(memoryUsageBefore); + } finally { + // release the flush and wait for its completion + releaseNativeWrite.countDown(); + kvTablet.setBeforeNativeWrite(null); + } flushAndWait(kvTablet, Long.MAX_VALUE); - assertThat(gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_ENTRY_COUNT)) - .isEqualTo(entryCountBefore); + + // the accounting returns to its previous value once the buffered entry is flushed assertThat(gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES)) .isEqualTo(memoryUsageBefore); } 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..c00403cc048 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 @@ -114,6 +114,7 @@ import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.stream.Collectors; import java.util.stream.IntStream; @@ -331,6 +332,11 @@ private KvTablet createKvTablet( schemaGetter, tableConf.getChangelogImage(), KvManager.getDefaultRateLimiter(), + null, + null, + new AtomicLong(), + new KvFlushScheduler(conf), + null, autoIncrementManager, clock, tableConf); @@ -351,6 +357,8 @@ private KvTablet createKvTablet( tableConf.getChangelogImage(), KvManager.getDefaultRateLimiter(), null, + null, + new AtomicLong(), kvFlushScheduler, null, autoIncrementManager, diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java index a277f02e11b..422c65194d9 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java @@ -26,6 +26,7 @@ import org.junit.jupiter.api.Test; import java.util.List; +import java.util.concurrent.atomic.AtomicLong; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; @@ -35,7 +36,8 @@ class KvPreWriteBufferTest { @Test void testIllegalLSN() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); bufferInsert(buffer, "key1", "value1", 1); bufferDelete(buffer, "key1", 3); @@ -54,7 +56,8 @@ void testIllegalLSN() { @Test void testWriteAndFlush() throws Exception { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); int elementCount = 0; // put a series of kv entries @@ -133,7 +136,8 @@ void testWriteAndFlush() throws Exception { @Test void testTruncate() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); int elementCount = 0; // put a series of kv entries @@ -186,7 +190,8 @@ void testTruncate() { @Test void testRowCount() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); int elementCount = 0; // put a series of kv entries @@ -230,7 +235,8 @@ void testRowCount() { @Test void testSplitPreparedFlushByRecordCount() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); // +key0(lsn 0), +key1(lsn 1), +key2(lsn 2), +key3(lsn 3), -key2(lsn 4) for (int i = 0; i < 4; i++) { bufferInsert(buffer, "key" + i, "value" + i, i); @@ -266,7 +272,8 @@ void testSplitPreparedFlushByRecordCount() { @Test void testSplitPreparedFlushByByteSize() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); // Entry payload sizes are 6, 5, and 5 bytes. bufferInsert(buffer, "a", "12345", 0); buffer.markWalBatchEnd(1); @@ -290,7 +297,8 @@ void testSplitPreparedFlushByByteSize() { @Test void testSplitPreparedFlushWithOversizedEntry() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); // The first entry is larger than the byte limit and must remain a non-empty singleton. bufferInsert(buffer, "a", "1234567890", 0); buffer.markWalBatchEnd(1); @@ -309,7 +317,8 @@ void testSplitPreparedFlushWithOversizedEntry() { @Test void testSplitPreparedFlushUsesFirstReachedLimit() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); // Entry payload sizes are 2, 2, 11, and 5 bytes. bufferInsert(buffer, "a", "x", 0); buffer.markWalBatchEnd(1); @@ -335,7 +344,8 @@ void testSplitPreparedFlushUsesFirstReachedLimit() { @Test void testCompletePrefixSegmentsAndAbortRest() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); // Entry payload sizes are 6, 5, 5, and 6 bytes. bufferInsert(buffer, "a", "12345", 0); buffer.markWalBatchEnd(1); @@ -382,7 +392,8 @@ private static void bufferDelete( @Test void testPrepareAndCompleteFlush() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); bufferInsert(buffer, "key1", "value1", 1); bufferInsert(buffer, "key2", "value2", 2); @@ -407,7 +418,8 @@ void testPrepareAndCompleteFlush() { @Test void testAbortPreparedFlush() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); bufferInsert(buffer, "key1", "value1", 1); bufferDelete(buffer, "key2", 2); @@ -427,7 +439,8 @@ void testAbortPreparedFlush() { @Test void testCannotTruncatePreparedFlush() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); bufferInsert(buffer, "key1", "value1", 1); bufferInsert(buffer, "key2", "value2", 2); @@ -440,7 +453,8 @@ void testCannotTruncatePreparedFlush() { @Test void testPendingFlushBytesTracking() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); // Initially zero assertThat(buffer.pendingFlushBytes()).isEqualTo(0); @@ -474,12 +488,11 @@ void testPendingFlushBytesTracking() { void testEstimatedMemoryUsage() { TabletServerMetricGroup metricGroup = new TabletServerMetricGroup(NOPMetricRegistry.INSTANCE, "fluss", "rack", "host", 0); - KvPreWriteBuffer buffer = new KvPreWriteBuffer(metricGroup); + AtomicLong sharedMemoryUsageBytes = new AtomicLong(); + KvPreWriteBuffer buffer = new KvPreWriteBuffer(metricGroup, sharedMemoryUsageBytes); assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); - assertThat(buffer.entryCount()).isEqualTo(0); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(0L); // +key1(10 bytes), +key2(11 bytes), -key3(4 bytes): 25 payload bytes in total bufferInsert(buffer, "key1", "value1", 1); @@ -489,55 +502,118 @@ void testEstimatedMemoryUsage() { // the estimation covers the payload bytes plus the per-entry object overhead assertThat(buffer.memoryUsageBytes()).isGreaterThan(payloadBytes); - assertThat(buffer.entryCount()).isEqualTo(3); - // the shared ledger mirrors the local accounting of the buffer - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()) - .isEqualTo(buffer.memoryUsageBytes()); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(3L); + // the shared counter mirrors the local accounting of the buffer + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(buffer.memoryUsageBytes()); // flushing all entries releases the whole accounted usage flushBuffer(buffer, Long.MAX_VALUE); assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); - assertThat(buffer.entryCount()).isEqualTo(0); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(0L); // truncating entries also releases their accounted usage bufferInsert(buffer, "key1", "value1", 4); assertThat(buffer.memoryUsageBytes()).isPositive(); buffer.truncateTo(4, TruncateReason.ERROR); assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); - assertThat(buffer.entryCount()).isEqualTo(0); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(0L); } @Test void testCloseReleasesAccounting() { TabletServerMetricGroup metricGroup = new TabletServerMetricGroup(NOPMetricRegistry.INSTANCE, "fluss", "rack", "host", 0); - KvPreWriteBuffer buffer = new KvPreWriteBuffer(metricGroup); + AtomicLong sharedMemoryUsageBytes = new AtomicLong(); + KvPreWriteBuffer buffer = new KvPreWriteBuffer(metricGroup, sharedMemoryUsageBytes); bufferInsert(buffer, "key1", "value1", 0); bufferInsert(buffer, "key2", "value2", 1); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isPositive(); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(2L); + assertThat(sharedMemoryUsageBytes.get()).isPositive(); - // closing releases the remaining accounting to the shared ledger exactly once + // closing releases the remaining accounting to the shared counter exactly once buffer.close(); assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); - assertThat(buffer.entryCount()).isEqualTo(0); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(0L); - // closing again is idempotent and must not over-release the ledger + // closing again is idempotent and must not over-release the shared counter buffer.close(); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().memoryUsageBytes()).isEqualTo(0L); - assertThat(metricGroup.kvPreWriteBufferMemoryLedger().entryCount()).isEqualTo(0L); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(0L); + } + + @Test + void testSameKeyVersionAccountingAcrossTruncateFlushClose() { + AtomicLong sharedMemoryUsageBytes = new AtomicLong(); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer( + TestingMetricGroups.TABLET_SERVER_METRICS, sharedMemoryUsageBytes); + KvPreWriteBuffer.Key key = toKey("k"); + + // reference: the accounting of a buffer holding just one version of the key + AtomicLong referenceSharedBytes = new AtomicLong(); + KvPreWriteBuffer reference = + new KvPreWriteBuffer( + TestingMetricGroups.TABLET_SERVER_METRICS, referenceSharedBytes); + bufferInsert(reference, "k", "v1", 1); + long singleVersionUsage = reference.memoryUsageBytes(); + + // v1 and v2 for the same key: only the latest version holds the key's map node + buffer.insert(key, "v1".getBytes(), 1); + buffer.insert(key, "v2".getBytes(), 2); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(buffer.memoryUsageBytes()); + + // truncating v2 reinstates v1 as the mapped version; the map-node accounting must be + // restored, leaving exactly the accounting of v1 alone instead of leaking it + buffer.truncateTo(2, TruncateReason.ERROR); + assertThat(buffer.memoryUsageBytes()).isEqualTo(singleVersionUsage); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(singleVersionUsage); + + // flushing the reinstated v1 brings both the local and the shared accounting back to + // zero with an empty buffer instead of going negative + flushBuffer(buffer, Long.MAX_VALUE); + assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(0L); + + // closing the empty buffer releases nothing extra and stays idempotent + buffer.close(); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(0L); + buffer.close(); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(0L); + } + + @Test + void testTruncatePartialCompletionStillReleasesAccounting() { + AtomicLong sharedMemoryUsageBytes = new AtomicLong(); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer( + TestingMetricGroups.TABLET_SERVER_METRICS, sharedMemoryUsageBytes); + + // k2 becomes PREPARED so a truncate past it fails midway, after the same-key k1 + // versions have already been removed + bufferInsert(buffer, "k2", "v", 1); + bufferInsert(buffer, "k1", "v1", 2); + bufferInsert(buffer, "k1", "v2", 3); + buffer.prepareFlush(2); + long usageBefore = sharedMemoryUsageBytes.get(); + + assertThatThrownBy(() -> buffer.truncateTo(1, TruncateReason.ERROR)) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("Cannot truncate prepared pre-write entry."); + + // the accounting of the already removed k1 versions, including the map node restored + // for the rolled-back version, must still be released to the shared counter: it equals + // the accounting of the remaining k2 entry and stays in sync with the local accounting + AtomicLong referenceSharedBytes = new AtomicLong(); + KvPreWriteBuffer reference = + new KvPreWriteBuffer( + TestingMetricGroups.TABLET_SERVER_METRICS, referenceSharedBytes); + bufferInsert(reference, "k2", "v", 1); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(reference.memoryUsageBytes()); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(buffer.memoryUsageBytes()); + assertThat(sharedMemoryUsageBytes.get()).isLessThan(usageBefore); } @Test void testCompleteFlushDetachesFlushedEntriesFromPreviousChain() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); KvPreWriteBuffer.Key key = toKey("k"); buffer.insert(key, "v1".getBytes(), 1); @@ -572,7 +648,8 @@ void testCompleteFlushDetachesFlushedEntriesFromPreviousChain() { @Test void testTruncateRollbackSemanticsAfterFlushDetachment() { - KvPreWriteBuffer buffer = new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS); + KvPreWriteBuffer buffer = + new KvPreWriteBuffer(TestingMetricGroups.TABLET_SERVER_METRICS, new AtomicLong()); KvPreWriteBuffer.Key key = toKey("k"); // v1@1, v2@2, v3@3; flush covers v1 and v2, v3 stays buffered diff --git a/website/docs/maintenance/observability/monitor-metrics.md b/website/docs/maintenance/observability/monitor-metrics.md index d7dbe387fbf..23075273de8 100644 --- a/website/docs/maintenance/observability/monitor-metrics.md +++ b/website/docs/maintenance/observability/monitor-metrics.md @@ -463,8 +463,8 @@ Some metrics might not be exposed when using other JVM implementations (e.g. IBM - tabletserver - - + tabletserver + - messagesInPerSecond The number of messages written per second to this server. Meter @@ -593,11 +593,6 @@ Some metrics might not be exposed when using other JVM implementations (e.g. IBM kvPreWriteBufferMemoryUsageBytes Estimated total memory usage of the KV pre-write buffers across all KV tablets in this server (in bytes), including the key/value payload bytes and the per-entry object overhead. It is an approximation for observability, not an exact measurement. Gauge - - - kvPreWriteBufferEntryCount - The number of entries buffered in the KV pre-write buffers across all KV tablets in this server. - Gauge kvWalMemoryPoolUsage