Skip to content

fix(scheduler): base EMALong.update() on the decayed value — fixes #301 - #308

Open
ImDanXie wants to merge 1 commit into
apache:mainfrom
ImDanXie:fix/scheduler-ema-decay-base
Open

ImDanXie wants to merge 1 commit into
apache:mainfrom
ImDanXie:fix/scheduler-ema-decay-base

Conversation

@ImDanXie

Copy link
Copy Markdown
Contributor

Summary

EMALong.get() applies time-based decay but never writes the decayed value back; EMALong.update() uses the raw, undecayed prev.ema as its base. The two methods therefore disagree about the current EMA value: after an idle period get() reports that congestion has subsided, yet the very next update() revives the stale value at 0.9× its old magnitude — and because update() refreshes lastTs, the revived value will not decay again for another decayDelay.

Failure scenario

With alpha=0.1, decay=0.9, decayDelay = maxBurstLatency = 5s (matching Batcher's defaults):

t=1s    update(20s) after a congestion spike      -> get() = 20000.0ms
t=61s   60s idle                                  -> get() = 60.9ms   (far below the 5s threshold)
t=61s   update(0s) after recovery                 -> get() = 18000.0ms

Batcher.submit() gates on emaQueueingTime.get() > maxBurstLatency, so after congestion has genuinely subsided the revived 18s value exceeds the 5s threshold and the batcher rejects requests with BackPressureException — whose callers drop messages outright. Every batch that does get admitted pushes the value back to 0.9×, producing an "admit one small batch, then drop for ~18s" cycle.

Fix

Extract the decay computation and use it in both get() and update(), so a fresh sample is blended onto the decayed value. (An alternative semantic — writing the decayed value back inside get() — would make a read modify the hot path's state; the extraction keeps reads side-effect free. Happy to switch if maintainers prefer the other variant.)

Test evidence

New regression test updateAfterIdleUsesDecayedBase: congestion spike (ema=20s), 60s idle, then update(0).

  • Control experiment: fails on the unfixed code (revives 18s), passes with the fix (~32ms).
  • Existing EMALong expectations unchanged (8/8).

Production validation

Deployed in a rolling 6-node cluster (wholesale lib/ replacement, one node at a time, zero rollbacks); cluster mesh, all 13 live client connections and cross-node delivery verified post-deploy. No back-pressure rejection storms observed in the rollout window.

Fixes #301

get() applies time-based decay but never writes it back, while update() used
the raw, undecayed prev.ema as its base. After an idle period the two methods
therefore disagreed: get() reported congestion had subsided, yet the next
update() revived the stale value at 0.9x its old magnitude and re-armed the
decay delay - which can pin the batcher in permanent back-pressure
("release a few batches, drop for tens of seconds" loop).

Extract the decay computation into a shared helper and use it in both get()
and update(), so a fresh sample is blended onto the decayed value.

Control experiment: updateAfterIdleUsesDecayedBase fails on the old code
(18s instead of ~32ms revived) and passes with the fix; existing EMALong
expectations unchanged (8/8).

Fixes apache#301
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.

[BUG] EMALong.update() uses the undecayed value as its base, resurrecting stale congestion after an idle period

1 participant