Lucas Bradstreet created KAFKA-21068:
----------------------------------------

             Summary: 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


[note: some edited/reduced LLM output below, human verified]

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.

The java client makes a bigger buffer to hold a record that is larger than 
`batch.size`. It then checks whether that bigger buffer is full, instead of 
checking whether the batch has reached the configured size target. 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:

```text
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
```

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.

## 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.

```diff
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];
```

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.

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