[
https://issues.apache.org/jira/browse/KAFKA-21068?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18114126#comment-18114126
]
sanghyeok An commented on KAFKA-21068:
--------------------------------------
Thanks for your answer.
self assigned.
> Produced records can wait for linger even after exceeding batch.size
> --------------------------------------------------------------------
>
> Key: KAFKA-21068
> URL: https://issues.apache.org/jira/browse/KAFKA-21068
> Project: Kafka
> Issue Type: Improvement
> Reporter: Lucas Bradstreet
> Assignee: sanghyeok An
> Priority: Major
>
> [note: some edited/reduced LLM output below, human verified and submitted]
> h4. Summary
> The java producer can delay a record that already exceeds batch.size because
> its buffer still has a few unused bytes.
>
> *Example:*
> With a 16 KiB batch size and 60-second linger, a single uncompressed 1 MiB
> record waits 60 seconds. Producing a 100-byte second record to the same
> partition after 20 seconds causes the first record to sent immediately.
> h4. Explanation
> A producer record substantially larger than `batch.size` can remain in an
> unready batch until `linger.ms` expires, even when the leader is known and
> there is no backoff, transaction fence, or buffer pressure.
> Another second record that fails to append to the same batch (due to lack of
> room) can make the first batch ready earlier and sendable. When a record
> exceeds batch.size, Java allocates a larger buffer to hold it. The problem is
> that Java then uses the larger buffer’s size as the new batching target. Even
> a few unused bytes can make Java keep waiting, although the record has
> already exceeded the configured batch size.
> The buffer includes a little extra space to be safe. Even a few unused bytes
> can therefore make Java keep waiting, although the batch is already much
> larger than the target.
> {*}{{*}}Compression is not required to trigger this issue.{{*}}{*} The
> simplest reproduction below uses `Compression.NONE`. Compression estimates
> can change whether the batch is considered full; the control results show
> both outcomes.
> In a deterministic uncompressed reproduction with `batch.size=16384` and a 1
> MiB value:
> {code:java}
> Configured batch target: 16,384 bytes
> Allocated capacity: 1,048,664 bytes
> Encoded/estimated size: 1,048,650 bytes
> Remaining allocation space: 14 bytes
> Batch isFull(): false
> Room for another 1 MiB record: false{code}
> With `linger.ms=1000`, this batch is unready at 0 and 999 ms, and becomes
> ready at 1000 ms without another record. A second 1 MiB append to the same
> partition makes it ready at time zero. With `linger.ms=0`, it is ready
> immediately.
> h4. Unit test of behavior (failure expected)
> Apply the following test-only diff to Kafka's existing producer
> `RecordAccumulatorTest`. The record is explicitly assigned to one partition
> to isolate the batching decision from sticky partition selection.
> The test appends one {*}{{*}}uncompressed 1 MiB value{{*}}{*} with
> `batch.size=16384`, `linger.ms=1000`, and an 8 MiB buffer pool. It checks
> that the batch already exceeds the target, then asserts that its leader is
> ready immediately. There is no second append or clock advance.
> {code:java}
> diff --git
> a/clients/src/test/java/org/apache/kafka/clients/producer/internals/RecordAccumulatorTest.java
>
> b/clients/src/test/java/org/apache/kafka/clients/producer/internals/RecordAccumulatorTest.java
> index 7890a5474b..075f3f9cf1 100644
> —
> a/clients/src/test/java/org/apache/kafka/clients/producer/internals/RecordAccumulatorTest.java
> +++
> b/clients/src/test/java/org/apache/kafka/clients/producer/internals/RecordAccumulatorTest.java
> @@ -276,6 +276,28 @@ public class RecordAccumulatorTest
> { testAppendLarge(Compression.NONE); }
> + @Test
> + public void testOversizedRecordIsReadyWithoutWaitingForLinger() throws
> Exception {
> + int batchSize = 16 * 1024;
> + int lingerMs = 1000;
> + RecordAccumulator accum = createTestRecordAccumulator(
> + batchSize, 8 * 1024 * 1024, Compression.NONE, lingerMs);
> + try
> { + long now = time.milliseconds(); +
> accum.append(topic, partition1, 0L, null, new byte[1024 * 1024], +
> Record.EMPTY_HEADERS, null, maxBlockTimeMs, now, cluster); + +
> assertEquals(1, accum.getDeque(tp1).size()); +
> assertTrue(accum.getDeque(tp1).peekFirst().estimatedSizeInBytes() >
> batchSize); + // No clock advance, second append, flush, or buffer
> pressure is needed. + assertEquals(Collections.singleton(node1),
> accum.ready(metadataCache, now).readyNodes, + "A batch already
> larger than batch.size should not wait for linger"); + }
> finally
> { + accum.abortIncompleteBatches(); + accum.close(); +
> }
> + }
> +
> private void testAppendLarge(Compression compression) throws Exception {
> int batchSize = 512;
> byte[] value = new byte[2 * batchSize];
> {code}
> The same diff is attached as `RecordAccumulatorTest.patch`. From the Kafka
> repository root, run:
> ```sh
> git apply RecordAccumulatorTest.patch
> ./gradlew :clients:test \
> --tests
> org.apache.kafka.clients.producer.internals.RecordAccumulatorTest.testOversizedRecordIsReadyWithoutWaitingForLinger
> ```
> This behavior reproduced against two real Kafka brokers with an uncompressed
> 1 MiB record, a 16 KiB batch size, and 60-second linger. Sending only that
> record resulted in an acknowledgement after 60.011 seconds. In a separate
> run, the record remained unsent for 20 seconds; appending a 100-byte second
> record then triggered delivery of the original record about 8 ms later,
> without a flush or close. With linger set to zero, the same 1 MiB record was
> acknowledged in 124 ms.
> Observed results
> |Codec / estimate state|Linger|Ready after first record at 0 ms?|Ready after
> next 1 MiB same-partition append at 0 ms?|Ready at linger deadline without
> another append?|
> |—|---:|—|—|—|
> |None|1000 ms|{*}{{*}}No{{*}}{*}|Yes|Yes; not ready at 999 ms|
> |None|0 ms|Yes|Not needed|Yes|
> |Gzip, initial estimate|1000 ms|Yes|Already ready|Already ready|
> |Zstd, initial estimate|1000 ms|Yes|Already ready|Already ready|
> |Gzip, explicit estimate 0.1|1000 ms|{*}{{*}}No{{*}}{*}|Yes|Yes; not ready at
> 999 ms|
> |Zstd, explicit estimate 0.1|1000 ms|{*}{{*}}No{{*}}{*}|Yes|Yes; not ready at
> 999 ms|
> The 0.1 compression estimate is injected with the existing
> `CompressionRatioEstimator.setEstimation` test helper. It represents a
> controlled already-learned state; no particular warmup workload is claimed to
> have learned that ratio. Initial gzip/zstd estimates produce 1,101,079
> estimated bytes, which exceed the allocated capacity and therefore do not
> reproduce the delayed-readiness case.
> h4. Historical evidence
> - {*}{{*}}0.10.2.0:{{*}}{*} the accumulator explicitly passed configured
> `batchSize` as the builder's write limit, separately from buffer allocation
> capacity. [Release
> source]([https://github.com/apache/kafka/blob/0.10.2.0/clients/src/main/java/org/apache/kafka/clients/producer/internals/RecordAccumulator.java])
> - {*}{{*}}24 March 2017:{{*}}{*} [KAFKA-4816 commit
> `5bd06f1d542`]([https://github.com/apache/kafka/commit/5bd06f1d542e6b588a1d402d059bc24690017d32])
> introduced the v2 format and conservative sizing, and changed the builder
> overload used by the accumulator to one that defaults its write limit to
> buffer capacity. The old final argument `this.batchSize` now represented a
> base offset rather than a write limit.
> - {*}{{*}}3 April 2017:{{*}}{*} [commit
> `f54b61909d5`]([https://github.com/apache/kafka/commit/f54b61909d525547d65123c02bbd36d92ccee5da])
> corrected that accidental base offset to `0L`; it did not restore the
> separate configured write limit.
> - {*}{{*}}0.11.0.0, released 28 June 2017:{{*}}{*} the released source
> contains the complete oversized-allocation/fullness/linger mechanism.
> [Release archive]([https://kafka.apache.org/community/downloads/#0.11.0.0])
> - Sticky partitioning arrived later through
> [KIP-480]([https://cwiki.apache.org/confluence/spaces/KAFKA/pages/120722025/KIP-480%2BSticky%2BPartitioner]),
> in Kafka 2.4.0. This behavior therefore predates sticky partitioning.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)