Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -590,6 +592,19 @@ private CompletableFuture<GCReply> 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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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<ROCoProcOutput> f = coProc.query(in, reader);
Expand Down Expand Up @@ -188,6 +195,39 @@ public void gcStopsOnQuota() {
assertEquals(r.getInspectedCount(), 3);
}

/**
* Regression: under step&gt;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<String> 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();
Expand Down