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 0219f3fe79..508b243d80 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,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 = 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 14a14714bc..b1db6ed343 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 61c5a2d5d7..4fd15ff056 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, @@ -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) { 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 18e922cf3f..9de3c56b84 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; @@ -87,7 +88,20 @@ * 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, + * 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<>(); @@ -99,15 +113,27 @@ public class KvPreWriteBuffer { private final Counter truncateAsDuplicatedCount; private final Counter truncateAsErrorCount; + // 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; // Accumulated byte size of entries not yet completed by a flush. private long pendingFlushBytes = 0; - public KvPreWriteBuffer(TabletServerMetricGroup serverMetricGroup) { + // 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; + + private boolean closed; + + public KvPreWriteBuffer( + TabletServerMetricGroup serverMetricGroup, AtomicLong memoryUsageBytesCounter) { truncateAsDuplicatedCount = serverMetricGroup.kvTruncateAsDuplicatedCount(); truncateAsErrorCount = serverMetricGroup.kvTruncateAsErrorCount(); + this.memoryUsageBytesCounter = memoryUsageBytesCounter; } /** @@ -181,8 +207,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); } /** @@ -201,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 */ @@ -211,39 +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; + } + } } - pendingFlushBytes -= entryBytes(entry.getKey(), entry.getValue()); - boolean removed = kvEntryMap.remove(entry.getKey(), 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. */ @@ -251,6 +296,16 @@ public long pendingFlushBytes() { return pendingFlushBytes; } + /** + * 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. * @@ -275,24 +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; - pendingFlushBytes -= entryBytes(entry.getKey(), entry.getValue()); - kvEntryMap.remove(entry.getKey(), 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()) { @@ -322,6 +391,58 @@ public void abortAllPrepared() { } } + /** + * 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) { + pendingFlushBytes += entryBytes(entry.getKey(), entry.getValue()); + long accountedBytes = entryAccountedBytes(entry, entry.previousEntry == null); + memoryUsageBytes += accountedBytes; + memoryUsageBytesCounter.addAndGet(accountedBytes); + } + + /** + * 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) { + pendingFlushBytes -= entryBytes(entry.getKey(), entry.getValue()); + boolean removedFromMap = kvEntryMap.remove(entry.getKey(), entry); + memoryUsageBytes -= entryAccountedBytes(entry, removedFromMap); + return removedFromMap; + } + + /** + * 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() { + if (closed) { + return; + } + closed = true; + memoryUsageBytesCounter.addAndGet(-memoryUsageBytes); + memoryUsageBytes = 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/metrics/group/TabletServerMetricGroup.java b/fluss-server/src/main/java/org/apache/fluss/server/metrics/group/TabletServerMetricGroup.java index dac5562c10..f3de1b34c0 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,9 @@ public class TabletServerMetricGroup extends AbstractMetricGroup { private volatile long sharedWriteBufferCapacity; + /** 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) { super(registry, new String[] {clusterId, hostname, NAME}, null); @@ -151,6 +154,10 @@ public TabletServerMetricGroup( // Register server-level RocksDB aggregated metrics registerServerRocksDBMetrics(); + + gauge( + MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES, + () -> kvPreWriteBufferMemoryUsageSupplier.getAsLong()); } /** @@ -215,6 +222,18 @@ public void setSharedWriteBufferMetrics(LongSupplier usageSupplier, long capacit this.sharedWriteBufferCapacity = capacity; } + /** + * 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 + */ + public void setKvPreWriteBufferMemoryUsageMetrics(LongSupplier usageSupplier) { + this.kvPreWriteBufferMemoryUsageSupplier = + checkNotNull(usageSupplier, "usageSupplier 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 e9c69e4fef..da87031f3e 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; @@ -282,6 +284,59 @@ void testSharedWriteBufferConfiguredThroughKvManagerCreateAndLoad() throws Excep .isEqualTo(capacity.getBytes()); } + @Test + void testPreWriteBufferServerLevelMetrics() throws Exception { + initTableBuckets(null); + TabletServerMetricGroup metricGroup = TestingMetricGroups.TABLET_SERVER_METRICS; + long memoryUsageBefore = + gaugeValue(metricGroup, MetricNames.KV_PRE_WRITE_BUFFER_MEMORY_USAGE_BYTES); + + KvTablet kvTablet = getOrCreateKv(tablePath1, null, tableBucket1); + // 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(); + } + }); + + 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); + + // 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); + } + @ParameterizedTest @MethodSource("partitionProvider") void testCreateKv(String partitionName) throws Exception { 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 568d9f9fa5..c00403cc04 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 45d983133a..422c65194d 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,13 +17,16 @@ 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; 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; @@ -33,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); @@ -52,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 @@ -131,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 @@ -184,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 @@ -228,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); @@ -264,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); @@ -288,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); @@ -307,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); @@ -333,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); @@ -380,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); @@ -405,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); @@ -425,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); @@ -438,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); @@ -468,9 +484,136 @@ void testPendingFlushBytesTracking() { assertThat(buffer.pendingFlushBytes()).isEqualTo(0); } + @Test + void testEstimatedMemoryUsage() { + TabletServerMetricGroup metricGroup = + new TabletServerMetricGroup(NOPMetricRegistry.INSTANCE, "fluss", "rack", "host", 0); + AtomicLong sharedMemoryUsageBytes = new AtomicLong(); + KvPreWriteBuffer buffer = new KvPreWriteBuffer(metricGroup, sharedMemoryUsageBytes); + + assertThat(buffer.memoryUsageBytes()).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); + 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); + // 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(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(sharedMemoryUsageBytes.get()).isEqualTo(0L); + } + + @Test + void testCloseReleasesAccounting() { + TabletServerMetricGroup metricGroup = + new TabletServerMetricGroup(NOPMetricRegistry.INSTANCE, "fluss", "rack", "host", 0); + AtomicLong sharedMemoryUsageBytes = new AtomicLong(); + KvPreWriteBuffer buffer = new KvPreWriteBuffer(metricGroup, sharedMemoryUsageBytes); + bufferInsert(buffer, "key1", "value1", 0); + bufferInsert(buffer, "key2", "value2", 1); + assertThat(sharedMemoryUsageBytes.get()).isPositive(); + + // closing releases the remaining accounting to the shared counter exactly once + buffer.close(); + assertThat(buffer.memoryUsageBytes()).isEqualTo(0L); + assertThat(sharedMemoryUsageBytes.get()).isEqualTo(0L); + + // closing again is idempotent and must not over-release the shared counter + buffer.close(); + 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); @@ -505,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 ee57382f52..23075273de 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 @@ -588,6 +588,11 @@ 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 kvWalMemoryPoolUsage