[Subscription] Consolidate poll queue fixes - #18771
Merged
jt2594838 merged 1 commit intoSep 30, 2026
Merged
Conversation
jt2594838
approved these changes
Sep 29, 2026
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
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
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.
Summary
This change consolidates four small fixes in the DataNode subscription polling path. Together they improve multi-topic poll deadline allocation, allow queue cleanup to proceed during long polls, preserve events encountered in transient non-pollable states, and recover prefetching after event-order disorder has drained.
Problem
pollV2held the queue read lock while waiting for data. Queue cleanup needs the write lock, so an empty or slow poll could delay cleanup.Implementation
pollV2read-lock scope to queue inspection, event state transitions, and in-flight registration. Waiting now happens outside the lock, with timer updates and interruption handling on each iteration.nackAndRequeue, which resets the event state and offers it back to the prefetch queue. Both polling paths use the helper.AtomicLongand clear the transient guard only after both prefetched and in-flight event counts reach zero.DataNodePipeMessagesconstants.Validation
mvn spotless:check -pl iotdb-core/datanode,iotdb-core/consensus,iotdb-core/node-commonsgit diff --checkThe targeted subscription test command reached DataNode compilation but was blocked by existing generated-source and dependency errors in unrelated
fillandmodeclasses (for example, missingIFillandAccumulatortypes). No compilation error was reported in the three changed source files before the unrelated generated-source failures.