You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
Fix consume with batch size 1 stopping after one message - #1633
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.
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.
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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
A simplified
consumewithConsumeOptions.batchSize(1)delivers the first message and then stalls. Every other batch size reads the whole stream.NatsMessageConsumercomputes the re-pull threshold as the batch size less the re-pull size:With
bm == 1the re-pull size floors at 1, so the threshold is 0. The check that schedules the next pull isand
pendingProcessedMessagesis 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
ConsumeBatchSizeTestsis 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.SimplificationTestspasses.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.