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)