Skip to content

Custom datanode backoff monitor 3.4 - #82

Draft
anthonyjdam wants to merge 13 commits into
hubspot-3.4from
custom-datanode-backoff-monitor-3.4
Draft

anthonyjdam wants to merge 13 commits into
hubspot-3.4from
custom-datanode-backoff-monitor-3.4

Conversation

@anthonyjdam

@anthonyjdam anthonyjdam commented Sep 23, 2026 •

Copy link
Copy Markdown

Load-aware HDFS decommission: DatanodeAdminAdaptiveBackoffMonitor

Human summary

Implements HBasePlanning#2591. Adds a new DatanodeAdminAdaptiveBackoffMonitor that uses hysteresis to ramp down DN decommissions aggressively when NN is unhealthy. The main NN load signal used is the average queue time for RPC requests (also has configurable window smoothing). The backoff monitor ships off by default, and there a bunch of knobs to fine tune it, which are documented in a follow up PR.

The problem

When we decommission a DataNode, HDFS has to re-replicate every block that node held onto the remaining nodes so the cluster returns to full replication before the host is removed. The NameNode drives this through a pluggable decommission monitor (dfs.namenode.decommission.monitor.class), and today we get to choose between exactly two upstream implementations — neither of which knows or cares how busy the NameNode is:

  • DatanodeAdminDefaultMonitor — scans a fixed number of blocks per interval while holding the FSNamesystem write lock for the whole tick. Simple, but the lock hold-time scales with the batch size and directly stalls foreground RPCs.
  • DatanodeAdminBackoffMonitor (HDFS-14854) — smarter: it paces itself off its own in-flight re-replication queue (pendingRep) and stops scheduling new work once that queue hits a fixed pendingRepLimit (default 10000). This bounds the lock pressure, but the cap is a static number set at boot.

Because the cap is fixed, it is always wrong in one of two directions:

  • Too aggressive when the NameNode is under load. A decommission kicked off during peak HBase traffic dumps up to 10000 in-flight re-replications into the same NameNode that's trying to serve reads and writes. Decommission work competes with foreground traffic for the FSNamesystem lock and RPC handlers, and client latency suffers — exactly when we can least afford it.
  • Too timid when the NameNode is idle. Overnight, or on a quiet cluster, that same 10000 cap needlessly throttles a decommission that could safely run far faster. Hosts sit in DECOMMISSION IN PROGRESS for hours longer than necessary, extending the window where we're running under-replicated and blocking hardware reclaim / rolls.

There's no single fixed value that's right for both. The cap needs to move with NameNode load: expand when the NameNode is healthy so decommission drains quickly, and contract when the NameNode is busy so foreground traffic is protected.

The solution

This PR adds a third monitor, DatanodeAdminAdaptiveBackoffMonitor, that subclasses the upstream backoff monitor and makes pendingRepLimit dynamic. Every tick it samples a NameNode-load signal and nudges the limit up or down with a closed-loop feedback controller:

  • NameNode healthy → ramp the limit up (additively) toward a configured ceiling → decommission runs as fast as the cluster can safely absorb.
  • NameNode busy → ramp the limit down (multiplicatively, faster than it ramps up) toward a floor → decommission yields NameNode capacity back to foreground traffic.
  • NameNode overloaded → an optional hard safety override slams the limit to the floor immediately.

Crucially, all of the actual block-scheduling machinery is inherited untouched from DatanodeAdminBackoffMonitor — the pendingRep queue, the lock cadence, node tracking and completion. The only thing this monitor changes is the value of the limit, recomputed once per tick. It is not a rewrite of the decommission loop; it's a governor bolted onto the existing one.

Observability quick reference

All metrics land on the NameNodeActivity JMX source and flow to VictoriaMetrics as collectd_hdfs_namenode_decommission_adaptive_*. The key one to watch is …_pending_limit — it drops when the controller backs off and climbs when it ramps up. See the admin docs (HBase → Cluster Ops → Adaptive DataNode Decommission Monitor) for the full metric list, tuning guide, and how to trigger a decommission.


Context for reviewer agents

Intent: Give HDFS decommission a load-aware pace on our HBase clusters. The pluggable NameNode decommission monitor (dfs.namenode.decommission.monitor.class) ships upstream as either DatanodeAdminDefaultMonitor (fixed blocks-per-interval, holds the FS write lock the whole tick) or DatanodeAdminBackoffMonitor (HDFS-14854, paces off its own in-flight pendingRep queue up to a fixed pendingRepLimit, default 10000). Neither reacts to NameNode health, so decommission-driven replication competes with foreground traffic when the NN is busy and runs needlessly slowly when it's idle. This PR adds a third monitor, DatanodeAdminAdaptiveBackoffMonitor, which makes that pending limit dynamic via a closed-loop feedback controller driven by a smoothed NameNode-load signal (average RPC queue time). It ships dark behind a flag (default off) — identical to the stock backoff monitor until enabled — and every knob is runtime-reconfigurable. Implements https://github.com/HubSpotEngineering/HBasePlanning/issues/2591. Base branch hubspot-3.3.6.

Changed files (12, +1681/−46):

  • ipc/metrics/RpcMetrics.java (+17) — expose the rolling RPC queue-time mean (hadoop-common).
  • hdfs/DFSConfigKeys.java (+60) — the config keys + defaults.
  • blockmanagement/DatanodeAdminAdaptiveBackoffMonitor.java (+518, new) — the monitor / controller.
  • blockmanagement/DatanodeAdminManager.java (+169) — reconfiguration refresh/get methods.
  • namenode/FSNamesystem.java (+45) — RPC load-signal accessors.
  • namenode/NameNode.java (+117) — RPC-server injection + reconfig wiring.
  • namenode/metrics/NameNodeMetrics.java (+74) — gauges + counters for adaptation observability.
  • resources/hdfs-default.xml (+121) — docs for every new key.
  • test/.../blockmanagement/TestDatanodeAdminAdaptiveBackoffMonitor.java (+301, new) — pure-unit tests.
  • test/.../hdfs/TestDecommissionWithAdaptiveBackoffMonitor.java (+55, new) — MiniDFSCluster e2e.
  • test/.../namenode/TestNameNodeReconfigure.java (+146) — runtime-reconfig tests.
  • test/.../tools/TestDFSAdmin.java (+104) — updates the reconfigurable-property assertion.

Diff shape:

  • RpcMetrics — adds getQueueMean() / getQueueSampleCount() (rpcQueueTime.lastStat().mean() / .numSamples()), mirroring the pre-existing getProcessingMean() / getProcessingSampleCount(). The rpcQueueTime MutableRate previously had no public getter.
  • FSNamesystem — a private volatile Server clientRpcServer (set from NameNode.initialize()), setClientRpcServer(Server), and two null-safe readers: getAvgRpcQueueTimeMs() (→ getQueueMean(), the primary signal) and getAvgRpcProcessingTimeMs() (→ getProcessingMean(), the optional processing-time override). Both return -1 when the server is unwired or no samples exist in the interval.
  • DatanodeAdminAdaptiveBackoffMonitor — extends DatanodeAdminBackoffMonitor; reuses all of the parent's tracking/backoff machinery and only changes how pendingRepLimit is chosen. Overrides processConf() (reads the keys, validateAndFixup(), derives the EWMA alpha) and run() (if enabled: setPendingRepLimit(computeAdaptivePendingLimit()), then super.run()). The control math is factored into two pure, @VisibleForTesting helpers — nextControllerLimit(current, smoothedSignal, forceMin) and smoothSignal(prevEma, sample) — so it's testable without the lock-holding run().
  • NameNode — injects namesystem.setClientRpcServer(rpcServer.getClientRpcServer()) in initialize(); adds the keys to the reconfigurableProperties set; adds an else if in reconfigurePropertyImpl → new reconfigureDecommissionAdaptiveMonitorParameters(...) handler that parses/validates each key and dispatches to DatanodeAdminManager.
  • NameNodeMetrics — gauges (DecommissionAdaptiveActive, …PendingLimit, …MinPendingLimit, …MaxPendingLimit, …RpcQueueTimeMs smoothed, …RawRpcQueueTimeMs) and counters (…RampUps, …RampDowns, …Holds, …ForceMins, …SignalUnavailable) on the NameNodeActivity source; the monitor updates them each tick via NameNode.getNameNodeMetrics() (best-effort, null-guarded). The action is classified once by nextControllerDecision(...) (returns a ControllerDecision of limit + ControllerAction), so the counter and the applied limit share one source of truth.
  • DatanodeAdminManager — a requireHubSpotMonitor(key) guard (instanceof + cast) and refresh*/get* per knob; adds ensurePositiveLong / ensureNonNegativeLong helpers (int knobs reuse ensurePositiveInt; the two override knobs use the existing ensureDisabledOrPositive).
  • hdfs-default.xml — one <property> per new key (required by TestHdfsConfigFields).
  • Tests — unit tests drive the controller step, the EWMA, config validation, and run() wiring; TestDecommissionWithAdaptiveBackoffMonitor runs the full TestDecommission suite with the monitor enabled; TestNameNodeReconfigure covers every knob + rejection when a non-adaptive monitor is active; TestDFSAdmin re-sorts the reconfigurable-property assertion (count 28 → 30).

How the controller works (per tick, only when enabled):

  1. Sample signal = FSNamesystem.getAvgRpcQueueTimeMs(). If < 0 (RPC metrics not available yet), fail open — return the current limit without advancing controller state.
  2. Seed the integral state from the current pendingRepLimit on the first tick (smooth enable, no step change).
  3. smoothed = smoothSignal(prevEma, signal) — optional cross-tick EWMA, alpha = tick / (window + tick).
  4. nextControllerLimit:
    • smoothed <= healthy → ramp up by ramp.up.step (capped at max);
    • smoothed >= busy → ramp down by ramp.down.step (floored at min);
    • in between → hold (the deadband is the hysteresis).
    • An optional hard override (busy.rpc.processing.time.ms) forces the floor immediately.

Config keys added (all under dfs.namenode.decommission.backoff.monitor., all reconfigurable at runtime):

key default meaning
adaptive.enabled false master switch; off ⇒ identical to stock backoff monitor
min.pending.limit 100 floor (busy); >0 so decommission never fully stalls
max.pending.limit inherits pending.limit (itself 10000) ceiling (healthy); when unset, defaults to the configured pending.limit, not a constant
healthy.rpc.queue.time.ms 1 avg RPC queue time (ms) at/below which the controller ramps the limit UP
busy.rpc.queue.time.ms 50 avg RPC queue time (ms) at/above which the controller ramps the limit DOWN; must be > healthy
ramp.up.step 500 blocks added to the limit per tick while healthy (slow up)
ramp.down.step 2000 blocks removed per tick while busy (fast down, AIMD-style)
signal.ema.window.ms 0 optional cross-tick EWMA window; 0 uses the RPC-metrics windowed mean directly
busy.rpc.processing.time.ms -1 optional hard override → floor; -1 disables

Development notes:

  • Same backoff engine, only the dial moves. The pendingRep queue, scheduling into neededReconstruction, blocksPerLock lock cadence, and node tracking/completion are 100% the parent DatanodeAdminBackoffMonitor's code, untouched. The only behavioral change is that pendingRepLimit is recomputed each tick. Reviewers should not expect changes to the block-scheduling loop.
  • Feedback (integral) controller, not absolute-target mapping. Modeled on HBase's FeedbackAdaptiveRateLimiter (and Janert's Feedback Control for Computer Systems): the limit is a control variable nudged each tick, not recomputed from a curve. The deadband between healthy and busy is the hysteresis and the small per-tick steps are the output-side smoothing, so adjacent ticks can never snap between floor and ceiling (which the earlier threshold-interpolation design did at both cliff edges). Ramp-down > ramp-up by default (fast to yield, slow to re-expand; AIMD).
  • Primary signal is average RPC queue time, deliberately not a block-queue count. getLowRedundancyBlocksCount() rises because the monitor schedules its own decommission work, so it would create a self-reinforcing loop. Average RPC queue time (how long calls wait before a handler picks them up) measures genuine foreground contention, is not inflated by our own scheduling, normalizes for service rate, and — critically — is a windowed mean already maintained by the RPC metrics system, so it survives the coarse (30s default) monitor tick where an instantaneous Server.getCallQueueLen() sample would just alias.
  • Known limits of the smoothing (call out for reviewers). getQueueMean() is the mean over the RPC metrics collection interval (configurable, commonly ~10s), which is shorter than the tick — so we read a rolling mean every tick rather than integrating the full inter-tick period; residual aliasing is greatly reduced, not zero. The in-monitor cross-tick EWMA (signal.ema.window.ms) exists to close that gap but is off by default (the metric is already a mean). DecayRpcScheduler's decayed averages were considered as an alternative smoothed source but not used, since they only exist when FairCallQueue is enabled; queue-time mean is always available.
  • Observability. Each tick the monitor publishes a full picture to NameNodeActivity, so the canary step is a dashboard/alert rather than a log-grep: DecommissionAdaptiveActive (1 only when actually pacing — distinguishes off / degraded / working), the effective limit with its Min/Max band, the RpcQueueTimeMs signal both smoothed and Raw, action-rate counters (RampUps/RampDowns/Holds/ForceMins — the "what is it doing and why", incl. the override firing), and SignalUnavailable (ticks the signal was -1, surfacing the sampling-window failure mode). All read 0 when the adaptive monitor isn't active; publishing is best-effort/null-guarded. NOTE: the metric values aren't unit-asserted (needs a running NameNode + live RPC samples); the classification that drives the counters is unit-tested via nextControllerDecision(...).action, the emit path is null-safe-covered by the run() tests, and a MetricsAsserts-based MiniDFSCluster test could be added if wanted.
  • Hard override vs. the control loop. busy.rpc.processing.time.ms is an orthogonal safety gate that slams the limit to the floor immediately (sustained high mean RPC processing time); it is not part of the proportional pacing and is disabled by default. (An earlier low-redundancy-block ceiling was dropped: lowRedundancy − pendingReconstruction subtracts across two different queues and during a real decommission stays ≈ our own contribution, so any cap value effectively pinned the monitor to the floor — a trap knob.)
  • Fail-open and self-disable. A -1 signal leaves the limit untouched; run() wraps the computation in a broad catch (Exception) (matching the parent run()) so adaptation can never break the decommission loop. validateAndFixup() clamps bad config at startup (min ≥ 1, max ≥ min, positive steps) and disables adaptation if busy <= healthy.
  • Ceiling inherits pending.limit. max.pending.limit defaults to the parent's resolved pending.limit (not a constant), and the reconfig reset path mirrors this via getConf(), so enabling adaptation never silently caps peak throughput below what the stock monitor already did.
  • Reconfig re-applies startup invariants. Every DatanodeAdminManager.refresh* validates the individual value, applies it, then re-runs the monitor's validateAndFixup(), so the cross-field invariants (max >= min, busy > healthy, and "disable adaptation on a broken deadband") hold identically at runtime and at boot. In particular, -reconfig adaptive.enabled=true with busy <= healthy still in place will not enable adaptation (it re-disables), and raising min above max lifts max to match rather than leaving min > max. Covered by TestNameNodeReconfigure.testReconfigureReappliesCrossFieldInvariants and two monitor-level unit tests.
  • Plumbing choice. Block-count signals were already reachable (the monitor holds blockManager); the only new wiring is the RPC-server reference into FSNamesystem, which avoids adding RPC-load accessors to the Namesystem interface (a namespace abstraction) just to satisfy this monitor. The base class types namesystem as Namesystem, so computeAdaptivePendingLimit() instanceof-guards the downcast rather than relying on the production wiring always passing an FSNamesystem: on anything else it fails open (holds the limit, i.e. behaves like the stock backoff monitor) and logs a single warning, instead of throwing a ClassCastException into run()'s catch every tick.

Anthony Dam and others added 13 commits September 8, 2026 16:29
Move all adaptive decommission knobs under the
dfs.namenode.decommission.backoff.monitor.adaptive.* namespace so they
are consistent with adaptive.enabled (previously only the enable flag
carried the "adaptive" segment, which made the tuning keys easy to
mis-set). Update hdfs-default.xml and the TestDFSAdmin sorted
reconfigurable-property assertion to match.

Also drop the per-tick smoothing diagnostic log and only emit the
pacing log when the effective limit actually changes; full per-tick
state remains on the NameNodeActivity metrics.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
@anthonyjdam
anthonyjdam marked this pull request as ready for review September 23, 2026 13:04
@anthonyjdam

Copy link
Copy Markdown
Author

After discussion, namenode stability is good right now and shipping this new monitor would incur a lot of effort and maintenance during future upgrades of Hadoop. Putting this on the back burner for now and marking as draft.

@anthonyjdam
anthonyjdam marked this pull request as draft September 23, 2026 16:22
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