[ 
https://issues.apache.org/jira/browse/KAFKA-21068?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

sanghyeok An reassigned KAFKA-21068:
------------------------------------

    Assignee: sanghyeok An

> 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)

Reply via email to