From 16c605bfbd8264170f80e009ee222b5030b2d28b Mon Sep 17 00:00:00 2001 From: Jiamei Xie Date: Thu, 12 Mar 2020 02:41:47 +0000 Subject: [PATCH] Free up compression buffers when splitting batches Method split will take up a lot of memory when the bigBatch is huge and compression is used. Call closeForRecordAppends() to free up resources like compression buffers. Change-Id: Iac6519fcc2e432330b8af2d9f68a8d4d4a07646b Signed-off-by: Jiamei Xie --- .../kafka/clients/producer/internals/ProducerBatch.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/ProducerBatch.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/ProducerBatch.java index 9323a61247c9a..cfd3a6794cc04 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/ProducerBatch.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/ProducerBatch.java @@ -268,14 +268,17 @@ public Deque split(int splitBatchSize) { // A newly created batch can always host the first message. if (!batch.tryAppendForSplit(record.timestamp(), record.key(), record.value(), record.headers(), thunk)) { batches.add(batch); + batch.closeForRecordAppends(); batch = createBatchOffAccumulatorForRecord(record, splitBatchSize); batch.tryAppendForSplit(record.timestamp(), record.key(), record.value(), record.headers(), thunk); } } // Close the last batch and add it to the batch list after split. - if (batch != null) + if (batch != null) { batches.add(batch); + batch.closeForRecordAppends(); + } produceFuture.set(ProduceResponse.INVALID_OFFSET, NO_TIMESTAMP, new RecordBatchTooLargeException()); produceFuture.done();