Skip to content

Fix consume with batch size 1 stopping after one message - #1633

Merged
scottf merged 2 commits into
mainfrom
consume-batch-size-one
Sep 28, 2026
Merged

scottf merged 2 commits into
mainfrom
consume-batch-size-one

Conversation

@scottf

@scottf scottf commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

A simplified consume with ConsumeOptions.batchSize(1) delivers the first message and then stalls. Every other batch size reads the whole stream.

NatsMessageConsumer computes the re-pull threshold as the batch size less the re-pull size:

int rePullMessages = Math.max(1, bm * consumeOpts.getThresholdPercent() / 100);
thresholdMessages = bm - rePullMessages;

With bm == 1 the re-pull size floors at 1, so the threshold is 0. The check that schedules the next pull is

if (pendingProcessedMessages < thresholdMessages) {
    processedHasCrossedThreshold = true;
}

and pendingProcessedMessages is itself clamped at 0, so it can never be less than 0. The consumer pulls once, delivers that message, and never pulls again. The byte threshold has the same shape for a batch bytes value small enough to floor the same way.

Fix: never let either threshold fall below 1.

New ConsumeBatchSizeTests is parameterized over batch sizes 1, 2, 3, 5 and 10 against a five message stream. Before the change, batch size 1 read 1 of 5 and the rest passed; after it, all five sizes read all five messages. SimplificationTests passes.

Found while testing a key value watcher built on the simplified consume; the same defect is in nats.java.v3 and is fixed there too.

🤖 Generated with Claude Code

Why this has gone unnoticed

A batch size of exactly 1 is the only value that reaches the bad threshold. The default batch is 500 with a 25% threshold, so the re-pull size is 125 and the threshold 375, far from the edge. A batch of 2 gives 2 - max(1, 0) = 1, which works, and every larger batch works for the same reason. Zero or a negative value is not the edge either - that falls back to the 500 default.

So it takes a caller deliberately setting batchSize(1), which is an unusual choice in real use, and the effect is silent rather than an error: one message arrives and the consumer simply goes quiet. It was found by sweeping batch sizes while testing something else, not by a report from the field.

@synadia-claude-reviewer

synadia-claude-reviewer Bot commented Sep 25, 2026 •

Copy link
Copy Markdown

Claude finished @scottf's task in 36s —— View job


Review

  • Read PR diff and description
  • Examine NatsMessageConsumer threshold logic and NatsMessageConsumerBase usage
  • Verify test coverage and helper infra
  • Post review

Correctness: The fix is sound. thresholdMessages/thresholdBytes are flooring correctly at 1, which matches the root cause described: with bm == 1, rePullMessages floors to 1 via Math.max(1, ...), so the un-patched thresholdMessages = bm - rePullMessages becomes 0 — a value pendingProcessedMessages (itself clamped to Math.max(0, ...)) can never drop below, so processedHasCrossedThreshold never flips true after the first pull. Flooring the threshold at 1 restores the intended "some processing occurred" signal. The same shape applies to thresholdBytes for small batchBytes, and that branch is fixed analogously (with the bb == 0 sentinel for "not tracking bytes" preserved correctly via Integer.MIN_VALUE).

I traced all read sites of thresholdMessages/thresholdBytes (statusAdjustPending, updateProcessed in NatsMessageConsumer.java) — the change is self-contained to this one constructor and doesn't affect other consumer types (e.g. NatsMessageConsumerBase/fetch-based consumers aren't touched).

Tests: ConsumeBatchSizeTests is a reasonable regression test — parameterized over batch sizes 1/2/3/5/10 against a 5-message stream, using existing JetStreamTestBase infra (jsServer, TestingStreamContainer, jsPublish) consistently with other tests in the suite.

No correctness, security, or performance concerns found. Nit (non-blocking): the inline comment above the threshold calc explains the fix well, but a comment restating the same rationale also exists in the new test's class Javadoc — fine as-is, just duplicated context.

@Jarema Jarema left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@scottf
scottf merged commit 3f1e0b4 into main Sep 28, 2026
5 checks passed
@scottf
scottf deleted the consume-batch-size-one branch September 28, 2026 18:59
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants