This is an automated email from the ASF dual-hosted git repository.
jsancio pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 05cbe688f75 KAFKA-20726: Don't reuse the BufferSupplier in kraft state
machine (#22678)
05cbe688f75 is described below
commit 05cbe688f75aadd1fb58b9d093201f7aa9e174c2
Author: Zhiyan Tang <[email protected]>
AuthorDate: Fri Jul 17 08:16:15 2026 -0500
KAFKA-20726: Don't reuse the BufferSupplier in kraft state machine (#22678)
KRaftControlRecordStateMachine reused a single BufferSupplier across all
reads. The supplier returned by BufferSupplier.create() caches
ByteBuffers by capacity without eviction, so the cached buffers lived as
long as the state machine and consumed memory unboundedly.
Create a fresh BufferSupplier per read in maybeLoadLog() and
maybeLoadSnapshot() instead, allowing it to be released once the read
completes.
Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
Reviewers: Kevin Wu <[email protected]>, Jonah Hooper
<[email protected]>, José Armando García Sancio <[email protected]>
---
.../org/apache/kafka/raft/KafkaRaftClient.java | 1 -
.../internals/KRaftControlRecordStateMachine.java | 47 ++++++------
.../KRaftControlRecordStateMachineTest.java | 83 +++++++++++++++++++++-
3 files changed, 105 insertions(+), 26 deletions(-)
diff --git a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java
b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java
index 824a4c4a8ba..44f6013e108 100644
--- a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java
+++ b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java
@@ -499,7 +499,6 @@ public final class KafkaRaftClient<T> implements
RaftClient<T> {
staticVoters,
log,
serde,
- BufferSupplier.create(),
MAX_BATCH_SIZE_BYTES,
logContext,
kafkaRaftMetrics,
diff --git
a/raft/src/main/java/org/apache/kafka/raft/internals/KRaftControlRecordStateMachine.java
b/raft/src/main/java/org/apache/kafka/raft/internals/KRaftControlRecordStateMachine.java
index 94bd8f52696..b5a9ccccdc2 100644
---
a/raft/src/main/java/org/apache/kafka/raft/internals/KRaftControlRecordStateMachine.java
+++
b/raft/src/main/java/org/apache/kafka/raft/internals/KRaftControlRecordStateMachine.java
@@ -55,7 +55,6 @@ public final class KRaftControlRecordStateMachine {
private final LogContext logContext;
private final RaftLog log;
private final RecordSerde<?> serde;
- private final BufferSupplier bufferSupplier;
private final Logger logger;
private final int maxBatchSizeBytes;
@@ -82,7 +81,6 @@ public final class KRaftControlRecordStateMachine {
* @param staticVoterSet the set of voter statically configured
* @param log the on disk topic partition
* @param serde the record decoder for data records
- * @param bufferSupplier the supplier of byte buffers
* @param maxBatchSizeBytes the maximum size of record batch
* @param logContext the log context
*/
@@ -90,7 +88,6 @@ public final class KRaftControlRecordStateMachine {
VoterSet staticVoterSet,
RaftLog log,
RecordSerde<?> serde,
- BufferSupplier bufferSupplier,
int maxBatchSizeBytes,
LogContext logContext,
KafkaRaftMetrics kafkaRaftMetrics,
@@ -100,7 +97,6 @@ public final class KRaftControlRecordStateMachine {
this.log = log;
this.voterSetHistory = new VoterSetHistory(staticVoterSet, logContext);
this.serde = serde;
- this.bufferSupplier = bufferSupplier;
this.maxBatchSizeBytes = maxBatchSizeBytes;
this.logger = logContext.logger(getClass());
this.kafkaRaftMetrics = kafkaRaftMetrics;
@@ -232,25 +228,27 @@ public final class KRaftControlRecordStateMachine {
}
private void maybeLoadLog() {
- while (log.endOffset().offset() > nextOffset) {
- LogFetchInfo info = log.read(
- nextOffset,
- Isolation.UNCOMMITTED,
- Integer.MAX_VALUE
- );
- try (RecordsIterator<?> iterator = new RecordsIterator<>(
- info.records,
- serde,
- bufferSupplier,
- maxBatchSizeBytes,
- true, // Validate batch CRC
- logContext
- )
- ) {
- while (iterator.hasNext()) {
- Batch<?> batch = iterator.next();
- handleBatch(batch, OptionalLong.empty());
- nextOffset = batch.lastOffset() + 1;
+ try (BufferSupplier bufferSupplier = BufferSupplier.create()) {
+ while (log.endOffset().offset() > nextOffset) {
+ LogFetchInfo info = log.read(
+ nextOffset,
+ Isolation.UNCOMMITTED,
+ Integer.MAX_VALUE
+ );
+ try (RecordsIterator<?> iterator = new RecordsIterator<>(
+ info.records,
+ serde,
+ bufferSupplier,
+ maxBatchSizeBytes,
+ true, // Validate batch CRC
+ logContext
+ )
+ ) {
+ while (iterator.hasNext()) {
+ Batch<?> batch = iterator.next();
+ handleBatch(batch, OptionalLong.empty());
+ nextOffset = batch.lastOffset() + 1;
+ }
}
}
}
@@ -268,7 +266,8 @@ public final class KRaftControlRecordStateMachine {
}
// Load the snapshot since the listener is at the start of the log
or the log doesn't have the next entry.
- try (SnapshotReader<?> reader = RecordsSnapshotReader.of(
+ try (BufferSupplier bufferSupplier = BufferSupplier.create();
+ SnapshotReader<?> reader = RecordsSnapshotReader.of(
rawSnapshot,
serde,
bufferSupplier,
diff --git
a/raft/src/test/java/org/apache/kafka/raft/internals/KRaftControlRecordStateMachineTest.java
b/raft/src/test/java/org/apache/kafka/raft/internals/KRaftControlRecordStateMachineTest.java
index 84ca54f206a..077b7da5b4d 100644
---
a/raft/src/test/java/org/apache/kafka/raft/internals/KRaftControlRecordStateMachineTest.java
+++
b/raft/src/test/java/org/apache/kafka/raft/internals/KRaftControlRecordStateMachineTest.java
@@ -34,12 +34,18 @@ import
org.apache.kafka.server.common.serialization.RecordSerde;
import org.apache.kafka.snapshot.RecordsSnapshotWriter;
import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
import org.mockito.Mockito;
+import java.util.ArrayList;
+import java.util.List;
import java.util.Optional;
import java.util.stream.IntStream;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.spy;
final class KRaftControlRecordStateMachineTest {
private static final RecordSerde<String> STRING_SERDE = new StringSerde();
@@ -58,7 +64,6 @@ final class KRaftControlRecordStateMachineTest {
staticVoterSet,
log,
STRING_SERDE,
- BufferSupplier.NO_CACHING,
1024,
new LogContext(),
raftMetrics,
@@ -461,4 +466,80 @@ final class KRaftControlRecordStateMachineTest {
assertEquals(Optional.of(voterSet),
partitionState.voterSetAtOffset(voterSetOffset));
assertEquals(4, getNumberOfVoters(metrics).metricValue());
}
+
+ @Test
+ void testBufferSupplierCreatedAndClosedOnLogRead() {
+ Metrics metrics = new Metrics();
+ KafkaRaftMetrics raftMetrics = new KafkaRaftMetrics(metrics, "raft");
+ ExternalKRaftMetrics externalMetrics =
Mockito.mock(ExternalKRaftMetrics.class);
+ MockLog log = buildLog();
+ VoterSet staticVoterSet =
VoterSetTest.voterSet(VoterSetTest.voterMap(IntStream.of(1, 2, 3), true));
+ int epoch = 1;
+
+ KRaftControlRecordStateMachine partitionState =
buildPartitionListener(log, staticVoterSet, raftMetrics, externalMetrics);
+
+ // Append a control record so that the log has something to read.
+ VoterSet voterSet =
VoterSetTest.voterSet(VoterSetTest.voterMap(IntStream.of(4, 5, 6), true));
+ log.appendAsLeader(
+ MemoryRecords.withVotersRecord(
+ log.endOffset().offset(),
+ 0,
+ epoch,
+ BufferSupplier.NO_CACHING.get(300),
+ voterSet.toVotersRecord((short) 0)
+ ),
+ epoch
+ );
+
+ BufferSupplier supplier = spy(new
BufferSupplier.GrowableBufferSupplier());
+ try (MockedStatic<BufferSupplier> bufferSupplierMock =
mockStatic(BufferSupplier.class)) {
+
bufferSupplierMock.when(BufferSupplier::create).thenReturn(supplier);
+
+ partitionState.updateState();
+
+ // The log read creates a BufferSupplier and closes it via
try-with-resources.
+ bufferSupplierMock.verify(BufferSupplier::create,
Mockito.times(1));
+ }
+
+ Mockito.verify(supplier).close();
+ }
+
+ @Test
+ void testBufferSupplierCreatedOnSnapshotRead() {
+ Metrics metrics = new Metrics();
+ KafkaRaftMetrics raftMetrics = new KafkaRaftMetrics(metrics, "raft");
+ ExternalKRaftMetrics externalMetrics =
Mockito.mock(ExternalKRaftMetrics.class);
+ MockLog log = buildLog();
+ VoterSet staticVoterSet =
VoterSetTest.voterSet(VoterSetTest.voterMap(IntStream.of(1, 2, 3), true));
+ int epoch = 1;
+
+ KRaftControlRecordStateMachine partitionState =
buildPartitionListener(log, staticVoterSet, raftMetrics, externalMetrics);
+
+ // Create a snapshot with control records so that the snapshot has
something to read.
+ KRaftVersion kraftVersion = KRaftVersion.KRAFT_VERSION_1;
+ VoterSet voterSet =
VoterSetTest.voterSet(VoterSetTest.voterMap(IntStream.of(4, 5, 6), true));
+ RecordsSnapshotWriter.Builder builder = new
RecordsSnapshotWriter.Builder()
+ .setRawSnapshotWriter(log.createNewSnapshotUnchecked(new
OffsetAndEpoch(10, epoch)).get())
+ .setKraftVersion(kraftVersion)
+ .setVoterSet(Optional.of(voterSet));
+ try (RecordsSnapshotWriter<?> writer = builder.build(STRING_SERDE)) {
+ writer.freeze();
+ }
+ log.truncateToLatestSnapshot();
+
+ List<BufferSupplier> suppliers = new ArrayList<>();
+ try (MockedStatic<BufferSupplier> bufferSupplierMock =
mockStatic(BufferSupplier.class)) {
+
bufferSupplierMock.when(BufferSupplier::create).thenAnswer(invocation -> {
+ BufferSupplier supplier = spy(new
BufferSupplier.GrowableBufferSupplier());
+ suppliers.add(supplier);
+ return supplier;
+ });
+
+ partitionState.updateState();
+ }
+
+ // Every BufferSupplier created during the read is closed via
try-with-resources.
+ assertFalse(suppliers.isEmpty());
+ suppliers.forEach(supplier -> Mockito.verify(supplier).close());
+ }
}