Skip to content

Add subscription stall diagnostics and reduce WAL prefetch overhead - #18808

Open
Caideyipi wants to merge 6 commits into
masterfrom
fix/subscription-stall-observability
Open

Caideyipi wants to merge 6 commits into
masterfrom
fix/subscription-stall-observability

Conversation

@Caideyipi

@Caideyipi Caideyipi commented Oct 9, 2026 •

Copy link
Copy Markdown
Collaborator

Description

Consensus subscription queues can stop delivering data while WAL lag remains positive. Poll/progress timestamps alone do not show gaps between data events, ongoing prefetch work, or allocator pressure. Periodic statistics pre-read WAL before pending delivery, per-entry pending cursor advances repeatedly inspect the same WAL footer, and Tablet memory rejection recreates the iterator and rescans the retained prefix. Add delivery/prefetch timing, allocator diagnostics, and DEBUG tracing, and remove these redundant WAL operations while preserving retry data and replay progress.

Behavior and metric semantics

  • Track data delivery with monotonic time when a data event is returned by poll. Empty polls, watermarks, and ACKs do not reset the data idle timer. Expose the last and maximum delivery interval and current data idle time in milliseconds.
  • Expose an ongoing prefetch round's duration, including queue lock acquisition, its maximum duration, and time since the last completed round. Timing gauge reads do not acquire the queue lock. Completed rounds lasting at least 10 seconds produce a diagnostic warning.
  • Export node subscription memory used, limit, overcommit bytes, and cumulative oversized-entry admissions directly from the shared allocator. These account for retained subscription allocations; temporary materialization and WAL/JVM allocations require separate JVM metrics.
  • Preserve master's per-queue memory handles, protected shares, and explicit oversized-entry rejection. Consensus queues do not admit an entry above their current per-queue maximum.
  • Retain the compatibility allocation API's single oversized-entry admission when the node budget is positive and otherwise empty. Its admission count and overcommit bytes are observable. Warn on the first admission and at most once every 30 seconds afterward. Rejected per-queue entries do not increment this admission count.
  • Include timing and allocator readings in queue core reports. Keep subscriptionMemoryUsedInBytes as the per-queue value and expose the node total separately as dataNodeSubscriptionMemoryUsedInBytes. Bind/unbind the new gauges across metric service restarts. Add matching English and Chinese messages.

Interpret delivery idle time together with queue active/initialized state, WAL lag, and polling activity, since it also grows while no data arrives or consumers do not poll.

Avoid redundant WAL work

  • Periodic statistics report walNextBuffered from the iterator's existing usable cache. Reading this field does not call hasNext(), read WAL, or advance replay. A false value does not imply exhaustion; interpret it alongside lag and worker activity.
  • Advance the local replay cursor as each pending request is processed, but call ProgressWALIterator.advanceTo() once at batch exit. This coalesces retained-file coverage checks for the processed prefix, including early exits for gaps or memory pressure. A seek generation change prevents stale batch work from fast-forwarding the iterator.
  • Keep a memory-rejected WAL request in the iterator's existing cache and retain its current reader and pending fragment state. After memory is available, retry the same request without reopening WAL or rescanning the prefix; successful materialization advances replay once.
  • Deserialize WAL fragments through duplicate ByteBuffer views so the original positions remain unchanged across retries. The views share payload bytes without making another payload copy. The retry retains one unmaterialized WAL request in the cache; iterator buffers and reader/JVM allocations remain outside the Tablet allocator gauges.
  • Preserve actual WAL readability checks in scheduling and gap recovery, as well as writer coverage checks when skipping files. Follower requests without a local search index remain eligible for replay.

Prefetch DEBUG tracing

Round start/end logs include a compact consumer-group/topic/region identifier, elapsed time, cursor changes, pending/WAL accepted-entry deltas, queue sizes, memory/admission blocking reasons, and the requested reschedule result/delay. Additional logs identify read-lock acquisition, WAL scan start, entry parsing/conversion and duration, estimated tablet bytes, and memory reservation rejection details. Subtask logs show whether a pending wakeup causes immediate re-enqueue despite the round's result.

To capture these logs, enable the following loggers in conf/logback-datanode.xml:

<logger name="org.apache.iotdb.db.subscription.broker.consensus.ConsensusPrefetchingQueue" level="DEBUG"/>
<logger name="org.apache.iotdb.db.subscription.task.subtask.ConsensusPrefetchSubtask" level="DEBUG"/>

Output is written to logs/log_datanode_debug.log with the worker thread name. DEBUG snapshots read existing state without probing the WAL iterator or constructing the queue's full core report.

Validation

  • Full 53-module English reactor test compilation and a full 53-module Chinese reactor test lifecycle run passed. English DataNode/upstream and Chinese full-reactor test runs enabled only the nine selected subscription test classes, with one Surefire fork.
  • Each locale ran 108 tests: 107 passed, no failures or errors, and the opt-in iterator performance benchmark was skipped. Classes: ProgressWALIteratorTest, ConsensusPrefetchingQueueTest, ConsensusPrefetchingQueueSeekTest, ConsensusPrefetchingQueueWalBackpressureTest, ConsensusLogToTabletConverterTest, ConsensusPrefetchingQueueDataNodeMemoryTest, SubscriptionMemoryManagerTest, ConsensusSubscriptionPrefetchingQueueMetricsTest, and SubscriptionMetricsRestartTest.
  • Regression checks verify that statistics never pre-read WAL, pending batches coalesce iterator fast-forward, memory backpressure advances only the consumed prefix, and cached follower requests survive local cursor advancement.
  • A real V3 WAL regression checks repeated memory rejection and ACK-driven recovery with the same iterator, ordered single acceptance of all requests, unchanged progress while waiting, bounded Tablet allocation, and final memory release. A two-fragment WAL deserialization regression checks that repeated attempts preserve every row, writer metadata, and original buffer positions.
  • Spotless, git diff --check, and English/Chinese message-key and placeholder parity passed.

Self-review

  • Reviewed concurrent timing reads/writes, allocator accounting, and subtask wakeup state; DEBUG logging occurs outside the subtask monitor.
  • Documented metric semantics, the compatibility allocator exception, and DEBUG activation.
  • Verified data gaps across watermarks/ACKs, timing during blocked prefetch, monotonic clock wrap, queue memory admission, and metric service restart/removal.
  • Included Apache license headers and matching localized messages.
  • Reviewed batched WAL fast-forward against pending gaps, memory backpressure, seek generations, and follower replay.
  • Verified memory retries preserve WAL payloads and reader state, and reviewed retry-cache invalidation on seek and cleanup.

Key changed classes

ConsensusPrefetchingQueue, ProgressWALIterator, ConsensusLogToTabletConverter, SubscriptionQueueTimeTracker, ConsensusPrefetchSubtask, ConsensusSubscriptionPrefetchingQueueMetrics, SubscriptionMetrics, and SubscriptionMemoryManager.

@codecov

codecov Bot commented Oct 9, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 3.81679% with 252 lines in your changes missing coverage. Please review.
✅ Project coverage is 45.81%. Comparing base (f16e5b3) to head (b7174d7).
⚠️ Report is 2 commits behind head on master.

Files with missing lines Patch % Lines
...on/broker/consensus/ConsensusPrefetchingQueue.java 0.00% 143 Missing ⚠️
...broker/consensus/SubscriptionQueueTimeTracker.java 0.00% 36 Missing ⚠️
.../ConsensusSubscriptionPrefetchingQueueMetrics.java 0.00% 36 Missing ⚠️
...db/db/subscription/metric/SubscriptionMetrics.java 0.00% 15 Missing ⚠️
...bscription/resource/SubscriptionMemoryManager.java 0.00% 13 Missing ⚠️
...ription/task/subtask/ConsensusPrefetchSubtask.java 0.00% 5 Missing ⚠️
...cription/broker/consensus/ProgressWALIterator.java 0.00% 3 Missing ⚠️
...roker/consensus/ConsensusLogToTabletConverter.java 0.00% 1 Missing ⚠️
Additional details and impacted files
@@             Coverage Diff              @@
##             master   #18808      +/-   ##
============================================
- Coverage     45.82%   45.81%   -0.02%     
  Complexity      712      712              
============================================
  Files          5498     5499       +1     
  Lines        397637   397857     +220     
  Branches      51765    51789      +24     
============================================
+ Hits         182235   182267      +32     
- Misses       215402   215590     +188     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@Caideyipi Caideyipi changed the title Add subscription delivery stall and memory overcommit diagnostics Add subscription stall diagnostics and reduce WAL prefetch overhead Oct 10, 2026
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.

1 participant