Environment
- BiFroMQ
4.0.0 / main @ eef5e3af
- Module:
base-scheduler
- File:
base-scheduler/src/main/java/org/apache/bifromq/basescheduler/EMALong.java:45-55 (update), :57-72 (get)
Summary
get() applies time-based decay (decay ^ seconds) but does not write the decayed value back to the state. update() uses the raw, undecayed prev.ema as its base. The two methods therefore disagree about the current value of the EMA: after a long idle period get() already reports that congestion has subsided, yet the very next update() revives the stale value at 90% of its old magnitude — and because update() refreshes lastTs, the revived value will not decay again for another decayDelay.
Code
public void update(long newValue) {
State prev = state.get();
long newEma = (prev.ema == 0L) ? newValue
: (long) Math.ceil(prev.ema * (1 - alpha) + newValue * alpha); // <-- prev.ema undecayed
...
}
public long get() {
...
double decayed = s.ema * Math.pow(decay, seconds);
return decayed < 1.0 ? 0L : Math.round(decayed); // <-- result not written back
}
Observed (before fix)
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; batcher would admit)
t=61s update(0s) after recovery -> get() = 18000.0ms
[ BUG ] at the same instant t=61s: get() = 60.9ms before update(), 18000.0ms after
The value at a single instant differs by a factor of ~300 depending on whether an update() happened. (Note: get() treats lastTs == 0 as uninitialised and returns the raw value.)
Impact
Batcher.submit() (Batcher.java:133) uses emaQueueingTime.get() > maxBurstLatency to apply back pressure. After congestion has genuinely subsided the revived 18s value exceeds the 5s threshold, so the batcher rejects requests with BackPressureException — whose callers drop messages outright. Worse, every batch that does get admitted calls update() again and pushes the value back to 0.9x, producing an "admit one small batch, then drop for ~18s" cycle until the historical value is gradually diluted below the threshold.
Suggested fix
Extract the decay computation and use it in both methods (semantics: update() blends the new sample onto the decayed value; get() and update() then agree on the current EMA). 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 and is the reason we propose this variant. Maintainer call on which semantic you prefer.
public void update(long newValue) {
long now = nowSupplier.get();
while (true) {
State prev = state.get();
long prevEma = decay(prev, now); // decayed base
long newEma = (prevEma == 0L) ? newValue
: (long) Math.ceil(prevEma * (1 - alpha) + newValue * alpha);
State next = new State(newEma, now);
if (state.compareAndSet(prev, next)) {
return;
}
}
}
public long get() {
return decay(state.get(), nowSupplier.get());
}
private long decay(State s, long now) { /* unchanged body of the old get() */ }
Observed (after fix)
[ OK ] at the same instant t=61s: get() = 60.9ms before update(), 54.8ms after
Test
New regression test: congestion spike (ema=20s) then 60s idle; asserts the next update(0) stays below 1s instead of reviving 18s. Existing EMALong expectations unchanged (8/8).
Production validation
A 6-node BifroMQ 4.0.0 cluster (standalone.sh, systemd, mTLS clientAuth=REQUIRE), fronted by a TCP load balancer with 6 backends. Live clients: 3 bridge replicas (MQTT5, $oshare subscribers) and 10 backend service instances (persistent sessions). RocksDB engine, default configuration otherwise.
The fix (plus other changes) was rolled across all six nodes one node at a time (wholesale lib/ replacement, then restart), with a per-node gate (service active, all four ports listening, cluster mesh re-established). Six of six nodes succeeded, zero rollbacks.
Post-deploy checks (all measured after the rollout):
| Check |
Result |
| Cluster mesh (inter-node 8898/8899 connections, per node) |
408–417 established on every node |
| Live client connections (3 bridge replicas + 10 backend services over mTLS) |
13/13 reconnected and stable |
| Cross-node delivery end-to-end |
publish on one node → bridge (connected elsewhere) consumed it → forwarded (counter 0 → 1) |
| Client-visible protocol smoke |
connect / subscribe / publish with QoS 1 all CONNACK 0 / PUBACK 0 |
No back-pressure rejection storms were observed in the rollout window; congestion/back-pressure indicators are being tracked against the pre-deploy baseline.
Environment
4.0.0/main @ eef5e3afbase-schedulerbase-scheduler/src/main/java/org/apache/bifromq/basescheduler/EMALong.java:45-55(update),:57-72(get)Summary
get()applies time-based decay (decay ^ seconds) but does not write the decayed value back to the state.update()uses the raw, undecayedprev.emaas its base. The two methods therefore disagree about the current value of the EMA: after a long idle periodget()already reports that congestion has subsided, yet the very nextupdate()revives the stale value at 90% of its old magnitude — and becauseupdate()refresheslastTs, the revived value will not decay again for anotherdecayDelay.Code
Observed (before fix)
With
alpha=0.1,decay=0.9,decayDelay = maxBurstLatency = 5s(matchingBatcher's defaults):The value at a single instant differs by a factor of ~300 depending on whether an
update()happened. (Note:get()treatslastTs == 0as uninitialised and returns the raw value.)Impact
Batcher.submit()(Batcher.java:133) usesemaQueueingTime.get() > maxBurstLatencyto apply back pressure. After congestion has genuinely subsided the revived 18s value exceeds the 5s threshold, so the batcher rejects requests withBackPressureException— whose callers drop messages outright. Worse, every batch that does get admitted callsupdate()again and pushes the value back to0.9x, producing an "admit one small batch, then drop for ~18s" cycle until the historical value is gradually diluted below the threshold.Suggested fix
Extract the decay computation and use it in both methods (semantics:
update()blends the new sample onto the decayed value;get()andupdate()then agree on the current EMA). An alternative semantic — writing the decayed value back insideget()— would make a read modify the hot path's state; the extraction keeps reads side-effect free and is the reason we propose this variant. Maintainer call on which semantic you prefer.Observed (after fix)
Test
New regression test: congestion spike (ema=20s) then 60s idle; asserts the next
update(0)stays below 1s instead of reviving 18s. Existing EMALong expectations unchanged (8/8).Production validation
A 6-node BifroMQ 4.0.0 cluster (
standalone.sh, systemd, mTLSclientAuth=REQUIRE), fronted by a TCP load balancer with 6 backends. Live clients: 3 bridge replicas (MQTT5,$osharesubscribers) and 10 backend service instances (persistent sessions). RocksDB engine, default configuration otherwise.The fix (plus other changes) was rolled across all six nodes one node at a time (wholesale
lib/replacement, then restart), with a per-node gate (service active, all four ports listening, cluster mesh re-established). Six of six nodes succeeded, zero rollbacks.Post-deploy checks (all measured after the rollout):
No back-pressure rejection storms were observed in the rollout window; congestion/back-pressure indicators are being tracked against the pre-deploy baseline.