Skip to content

[Subscription] Consolidate poll queue fixes - #18771

Merged
jt2594838 merged 1 commit into
apache:masterfrom
Caideyipi:codex/subscription-four-20260929
Sep 30, 2026
Merged

jt2594838 merged 1 commit into
apache:masterfrom
Caideyipi:codex/subscription-four-20260929

Conversation

@Caideyipi

@Caideyipi Caideyipi commented Sep 29, 2026 •

Copy link
Copy Markdown
Collaborator

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

  • The per-topic poll deadline was divided only among topics after the current topic. The first topic could consume too much of the receiver deadline, leaving later topics little or no time to poll.
  • pollV2 held the queue read lock while waiting for data. Queue cleanup needs the write lock, so an empty or slow poll could delay cleanup.
  • A prefetched event that was temporarily non-pollable was nacked after being removed from the queue, but was not restored to the queue. This could silently drop the event.
  • A single disorder event permanently disabled prefetching for the queue, even after all events from the affected sequence had been consumed.

Implementation

  • Capture the current topic count before decrementing the remaining-topic counter and use it when dividing the receiver's remaining deadline.
  • Limit the pollV2 read-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.
  • Centralize the non-pollable recovery path in nackAndRequeue, which resets the event state and offers it back to the prefetch queue. Both polling paths use the helper.
  • Replace the disorder counter with an AtomicLong and clear the transient guard only after both prefetched and in-flight event counts reach zero.
  • Route the added polling warnings through the existing localized DataNodePipeMessages constants.

Validation

  • mvn spotless:check -pl iotdb-core/datanode,iotdb-core/consensus,iotdb-core/node-commons
  • git diff --check
  • Maven Checkstyle completed with zero violations during the targeted verification run.

The targeted subscription test command reached DataNode compilation but was blocked by existing generated-source and dependency errors in unrelated fill and mode classes (for example, missing IFill and Accumulator types). No compilation error was reported in the three changed source files before the unrelated generated-source failures.

@jt2594838
jt2594838 merged commit 8f1bed8 into apache:master Sep 30, 2026
37 of 39 checks passed
@jt2594838
jt2594838 deleted the codex/subscription-four-20260929 branch September 30, 2026 02:37
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