From f9500bd710708104e163a2fe6a0d3aaf3123d9d8 Mon Sep 17 00:00:00 2001 From: ImDanXie Date: Sat, 26 Sep 2026 01:58:35 +0800 Subject: [PATCH] fix(dist): rotate the GC sampling phase so stepping cannot lock the sweep cursor Under step>1 the co-proc inspects every step-th key. The cleaner resumes each session at the key where the previous one stopped, and on wrap that key is the session's own start key - so the cursor is pinned and the same modulo class is inspected in every session while the remaining routes are never scanned (measured: 50% of routes permanently skipped at step 2, stable across 400/900 rounds). Dead routes in the skipped classes are never collected and the route table grows without bound. Advance a per-session phase offset (0..step-1) before the first inspection, so consecutive sessions visit every residue class. No protocol change: the phase is derived from a per-range session counter inside the co-proc. Fixes #300 --- .../bifromq/dist/worker/DistWorkerCoProc.java | 15 +++++++ .../dist/worker/DistWorkerCoProcGCTest.java | 40 +++++++++++++++++++ 2 files changed, 55 insertions(+) diff --git a/bifromq-dist/bifromq-dist-worker/src/main/java/org/apache/bifromq/dist/worker/DistWorkerCoProc.java b/bifromq-dist/bifromq-dist-worker/src/main/java/org/apache/bifromq/dist/worker/DistWorkerCoProc.java index d49176e90..85d6dd2b7 100644 --- a/bifromq-dist/bifromq-dist-worker/src/main/java/org/apache/bifromq/dist/worker/DistWorkerCoProc.java +++ b/bifromq-dist/bifromq-dist-worker/src/main/java/org/apache/bifromq/dist/worker/DistWorkerCoProc.java @@ -108,6 +108,8 @@ class DistWorkerCoProc implements IKVRangeCoProc { private final ITenantsStats tenantsState; private final IDeliverExecutorGroup deliverExecutorGroup; private final ISubscriptionCleaner subscriptionChecker; + /** Per-range GC session counter used to rotate the sampling phase under step>1. */ + private final java.util.concurrent.atomic.AtomicInteger gcSessionSeq = new java.util.concurrent.atomic.AtomicInteger(); private transient Fact fact; private transient Boundary boundary; @@ -590,6 +592,19 @@ private CompletableFuture gc(GCRequest request, IKVRangeReader reader) .build()); } + // Rotate the sampling phase across sessions. The cleaner resumes each session at the + // key where the previous one stopped; on wrap that key is the session's own start key, so with a + // fixed step the same modulo class was inspected forever while the rest of the route table was + // never scanned (measured: 50% of routes permanently skipped at step 2). Advancing a per-session + // phase offset guarantees every residue class is visited over consecutive sessions. + int phase = stepUsed > 1 ? Math.floorMod(gcSessionSeq.getAndIncrement(), stepUsed) : 0; + while (phase-- > 0) { + itr.next(); + if (!itr.isValid()) { + break; // tail reached; the scan loop below wraps to the first key + } + } + AtomicInteger inspectedCount = new AtomicInteger(); AtomicBoolean wrapped = new AtomicBoolean(false); ByteString sessionStartKey = null; diff --git a/bifromq-dist/bifromq-dist-worker/src/test/java/org/apache/bifromq/dist/worker/DistWorkerCoProcGCTest.java b/bifromq-dist/bifromq-dist-worker/src/test/java/org/apache/bifromq/dist/worker/DistWorkerCoProcGCTest.java index 5c2f9e8f0..7608c3bd7 100644 --- a/bifromq-dist/bifromq-dist-worker/src/test/java/org/apache/bifromq/dist/worker/DistWorkerCoProcGCTest.java +++ b/bifromq-dist/bifromq-dist-worker/src/test/java/org/apache/bifromq/dist/worker/DistWorkerCoProcGCTest.java @@ -103,6 +103,10 @@ public void setUp() { } private ROCoProcOutput gc(long reqId, Integer stepHint, Integer scanQuota) { + return gcFrom(reqId, stepHint, scanQuota, null); + } + + private ROCoProcOutput gcFrom(long reqId, Integer stepHint, Integer scanQuota, ByteString startKey) { when(reader.iterator()).thenReturn(new FakeIterator(currentData)); GCRequest.Builder req = GCRequest.newBuilder().setReqId(reqId); if (stepHint != null) { @@ -111,6 +115,9 @@ private ROCoProcOutput gc(long reqId, Integer stepHint, Integer scanQuota) { if (scanQuota != null) { req.setScanQuota(scanQuota); } + if (startKey != null) { + req.setStartKey(startKey); + } ROCoProcInput in = ROCoProcInput.newBuilder() .setDistService(DistServiceROCoProcInput.newBuilder().setGc(req.build()).build()).build(); CompletableFuture f = coProc.query(in, reader); @@ -188,6 +195,39 @@ public void gcStopsOnQuota() { assertEquals(r.getInspectedCount(), 3); } + /** + * Regression: under step>1 the cleaner's cursor stays pinned at the same start key whenever a + * session wraps, so without phase rotation the same modulo class of the route table would be inspected + * in every session and the other half would never be scanned (measured: 50% stable gap at step 2). + * Consecutive sessions must rotate the sampling phase until every route has been inspected. + */ + @Test + public void gcSessionsRotatePhaseSoAllKeysGetInspected() { + currentData.clear(); + int total = 12; + for (int i = 0; i < total; i++) { + currentData.add(new KV(normalKey("t" + String.format("%02d", i), "p/" + i, "1\0i" + i + "\0d1"), + toByteString(1L))); + } + java.util.Set inspectedTenants = new java.util.HashSet<>(); + when(subscriptionChecker.sweep(any(Integer.class), any(CheckRequest.class))) + .thenAnswer(inv -> { + CheckRequest req = inv.getArgument(1); + inspectedTenants.add(req.getTenantId()); + return CompletableFuture.completedFuture(new ISubscriptionCleaner.GCStats(req.getMatchInfoCount(), 0)); + }); + + ByteString startKey = null; + long reqId = 3000L; + for (int session = 0; session < 8 && inspectedTenants.size() < total; session++) { + GCReply r = gcFrom(reqId++, 2, 100, startKey).getDistService().getGc(); + startKey = r.hasNextStartKey() ? r.getNextStartKey() : null; + } + + assertEquals(inspectedTenants.size(), total, + "every route must be inspected within a few sessions under step=2"); + } + @Test public void gcWrapAndStopOnSessionStartKey() { currentData.clear();