Skip to content

Bound subscription memory borrowing and report admission failures - #18812

Merged
jt2594838 merged 2 commits into
masterfrom
fix/subscription-fair-share-admission
Oct 10, 2026
Merged

jt2594838 merged 2 commits into
masterfrom
fix/subscription-fair-share-admission

Conversation

@Caideyipi

@Caideyipi Caideyipi commented Oct 9, 2026 •

Copy link
Copy Markdown
Collaborator

Description

In a mixed large-field and numeric subscription workload, one queue retained most of a DataNode's subscription materialization budget while other queues also rejected real-time entries. The old warning called every rejection "queue full", even when capacity remained. The reported run (case 1392, result 27524) consumed 18,919,210 of 24,693,910 written rows before its consumer timed out. The source counts matched the writes. This PR bounds how much materialized-tablet memory one queue can retain and makes admission failures observable; it does not claim that this allocation change alone explains or resolves the entire observed delivery stall.

Admission and recovery

  • Protect half of each active queue's equal share of the DataNode materialization budget. A queue can borrow idle capacity above that floor but cannot consume an existing peer's protected floor. For example, with two queues and a 100-byte budget, one may retain up to 75 bytes while the other keeps at least 25 bytes available. Borrowed allocations cannot be revoked; a newly activated queue may wait for delivery and ACK before its floor becomes available.
  • Preserve already accepted pending requests during ordinary backpressure. Cleanup and seek still clear stale pending requests, and rejected real-time requests remain eligible for WAL replay.
  • Return a critical poll error containing SUBSCRIPTION_OVERSIZED_ENTRY only when one materialized WAL entry exceeds the maximum a queue can hold under the current protected floors. A larger-than-equal-share entry can use idle borrowable capacity. Do not advance an oversized entry's WAL or commit progress. This is an intentional behavior change from the empty-budget oversized-entry exception; operators must increase the materialization budget, reduce concurrent queues, or narrow the topic to consume such an entry.

Diagnostics and configuration

  • Report a stable reasonCode for real-time admission rejection, distinguishing subscription quota/limit, consensus request memory, writer backlog, queue capacity, and lifecycle states. Add queue-level memory used/equal share/maximum and total plus per-reason rejection counters to metrics. Rate-limit WARN counts separately by reason, and include the latest reason in each queue core report. Capture offer rejection reasons per writer thread so concurrent offers cannot mislabel one another. A transient real-time rejection is still eligible for WAL replay and does not become a consumer error for every rejected offer.
  • Recommend the single starting value subscription_cache_memory_usage_percentage=0.2 in the distributed config template, including large-field and large-batch workloads. This controls serialized poll-response caching, not the materialized-tablet budget. Size the ninth chunk_timeseriesmeta_free_memory_proportion share from the largest materialized entry, active queue count, and desired buffering depth.

Validation and limits

  • 58 focused tests passed: 56 DataNode and 2 consensus tests, including protected borrowing, new-queue drain, mixed-size queues, pending preservation, WAL recovery, oversized-entry errors, lifecycle cleanup, and concurrent rejection reporting.
  • The full English reactor test build and Chinese reactor test-compile -DskipTests -P with-zh-locale passed on the final changes. Spotless and git diff --check passed.
  • The reported 3c1d/2-writer workload has not been rerun on this change. The 960 MB retained-tablet snapshot in result 27524 was taken after the consumer exited, so it does not by itself prove the cause of the earlier 67-second stop. Re-run that workload while collecting rejection codes, memory pools, prefetch-task durations, and ACK/GC timing.
  • Queue isolation applies to the subscription tablet materialization budget. Consensus request memory and prefetch workers are still shared. A peer entry larger than its protected floor may wait for borrowed memory to drain after ACK; this policy guarantees the floor for already active queues, not immediate recovery of borrowed bytes or all entry sizes.

Related PR

#18808 adds delivery-stall and node-memory diagnostics and retains the previous single oversized-entry exception. These PRs overlap in ConsensusPrefetchingQueue, metrics, memory manager, and localized messages. Their metric and oversized-entry semantics need reconciliation before merging both.


This PR has:

  • been self-reviewed for concurrent reads and writes.
  • updated configuration guidance and Javadocs for the changed behavior.
  • added focused unit tests for new paths.
  • been tested in a test IoTDB cluster with the reported workload.

Key changed classes

SubscriptionMemoryManager, ConsensusPrefetchingQueue, SubscriptionQueueRegistry, and ConsensusSubscriptionPrefetchingQueueMetrics.

@Caideyipi
Caideyipi force-pushed the fix/subscription-fair-share-admission branch from 57a1f84 to 2bfe91b Compare October 9, 2026 08:55
@Caideyipi
Caideyipi force-pushed the fix/subscription-fair-share-admission branch from 2bfe91b to 1dc912f Compare October 9, 2026 09:01
@Caideyipi Caideyipi changed the title Isolate consensus subscription memory by queue and report admission failures Bound subscription memory borrowing and report admission failures Oct 9, 2026
@codecov

codecov Bot commented Oct 9, 2026

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 7.24234% with 333 lines in your changes missing coverage. Please review.
✅ Project coverage is 45.55%. Comparing base (eb0e3f1) to head (80a8ca8).
⚠️ Report is 1 commits behind head on master.

Files with missing lines Patch % Lines
...bscription/resource/SubscriptionMemoryManager.java 0.00% 148 Missing ⚠️
...on/broker/consensus/ConsensusPrefetchingQueue.java 0.00% 120 Missing ⚠️
.../ConsensusSubscriptionPrefetchingQueueMetrics.java 0.00% 51 Missing ⚠️
...us/iot/subscription/SubscriptionQueueRegistry.java 27.77% 13 Missing ⚠️
...subscription/SubscriptionQueueRejectionReason.java 92.85% 1 Missing ⚠️
Additional details and impacted files
@@             Coverage Diff              @@
##             master   #18812      +/-   ##
============================================
- Coverage     45.59%   45.55%   -0.04%     
  Complexity      712      712              
============================================
  Files          5496     5497       +1     
  Lines        396800   397113     +313     
  Branches      51621    51648      +27     
============================================
- Hits         180929   180921       -8     
- Misses       215871   216192     +321     

☔ 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.

Comment on lines +153 to +156
if (!handle.active || handle.closed || !handles.containsKey(handle.id)) {
return AllocationResult.rejected(
AllocationRejectionReason.MEMORY_LIMIT, getTotalMemorySizeInBytes(), 0L);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Why is this called MEMORY_LIMIT?

@jt2594838
jt2594838 merged commit 21d53dd into master Oct 10, 2026
37 of 40 checks passed
@jt2594838
jt2594838 deleted the fix/subscription-fair-share-admission branch October 10, 2026 03:17
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