From bd7ce2aa2e1f39f2d3c923b3f29941c4393fc927 Mon Sep 17 00:00:00 2001
From: Doug Lea
Date: Fri, 31 Jul 2026 11:09:54 -0400
Subject: [PATCH] Test more guesses about FMR slowdowns
---
.../java/util/concurrent/ForkJoinPool.java | 186 ++++++++++--------
1 file changed, 106 insertions(+), 80 deletions(-)
diff --git a/src/java.base/share/classes/java/util/concurrent/ForkJoinPool.java b/src/java.base/share/classes/java/util/concurrent/ForkJoinPool.java
index 84e47cfa1d0..55902d9de76 100644
--- a/src/java.base/share/classes/java/util/concurrent/ForkJoinPool.java
+++ b/src/java.base/share/classes/java/util/concurrent/ForkJoinPool.java
@@ -1081,6 +1081,7 @@ public class ForkJoinPool extends AbstractExecutorService
static final long CLEANED = 1L << 2; // stopped and queues cleared
static final long TERMINATED = 1L << 3; // only set if STOP also set
static final long RS_LOCK = 1L << 4; // lowest seqlock bit
+ static final long RS_SEQ = 1L << 5; // version bit
// spin/sleep limits for runState locking and elsewhere
static final int LOCK_SPINS = 1 << 7; // max onSpinWaits for locks
@@ -1162,6 +1163,11 @@ public class ForkJoinPool extends AbstractExecutorService
return (ForkJoinTask>)U.getReferenceAcquire(a, slotOffset(idx));
}
+ @jdk.internal.vm.annotation.DontInline
+ static void doExec(ForkJoinTask> t) {
+ t.doExec();
+ }
+
// Nested classes
/**
@@ -1289,7 +1295,7 @@ public class ForkJoinPool extends AbstractExecutorService
a[j] = task;
U.fullFence();
if (room == 0) // resize
- growArray(a, cap, s);
+ growArray(a, cap, s, pool);
else if (room != m && (pj = m & (s - 1)) >= 0 && pj < cap &&
a[pj] != null)
pool = null;
@@ -1309,7 +1315,8 @@ public class ForkJoinPool extends AbstractExecutorService
* @param cap old array capacity
* @param s current top
*/
- private void growArray(ForkJoinTask>[] a, int cap, int s) {
+ private void growArray(ForkJoinTask>[] a, int cap, int s,
+ ForkJoinPool pool) {
int newCap = (cap < (1 << 16)) ? cap << 2 : cap << 1;
if (a != null && a.length == cap && cap > 0 && newCap > 0) {
ForkJoinTask>[] newArray = null;
@@ -1325,6 +1332,8 @@ public class ForkJoinPool extends AbstractExecutorService
break; // lost to pollers
}
U.putReferenceRelease(this, ARRAY, newArray);
+ if (pool != null) // ensure rescans
+ pool.advanceRunState();
}
}
}
@@ -1420,7 +1429,7 @@ public class ForkJoinPool extends AbstractExecutorService
int fifo = config & FIFO;
++nsteals;
while (task != null) {
- task.doExec();
+ doExec(task);
task = null;
int p = top, cap; ForkJoinTask>[] a;
if ((a = array) == null || (cap = a.length) <= 0)
@@ -1475,7 +1484,7 @@ public class ForkJoinPool extends AbstractExecutorService
if (!internal)
unlockPhase();
if (taken)
- task.doExec();
+ doExec(task);
break;
}
}
@@ -1522,7 +1531,7 @@ public class ForkJoinPool extends AbstractExecutorService
if (!casSlotNull(a, j, t))
break;
top = s;
- t.doExec();
+ doExec(t);;
if (limit != 0 && --limit == 0)
break;
}
@@ -1560,7 +1569,7 @@ public class ForkJoinPool extends AbstractExecutorService
unlockPhase();
if (!taken)
break;
- t.doExec();
+ doExec(t);;
if (limit != 0 && --limit == 0)
break;
}
@@ -1596,7 +1605,7 @@ public class ForkJoinPool extends AbstractExecutorService
else if (base == b && a[j] == t && casSlotNull(a, j, t)) {
base = b + 1;
U.storeFence();
- t.doExec();
+ doExec(t);;
}
}
}
@@ -1742,6 +1751,9 @@ public class ForkJoinPool extends AbstractExecutorService
private void unlockRunState() { // increment lock bit
U.getAndAddLong(this, RUNSTATE, RS_LOCK);
}
+ final void advanceRunState() { // increment version bit
+ U.getAndAddLong(this, RUNSTATE, RS_SEQ);
+ }
private long lockRunState() { // lock and return current state
long s, u; // locked when RS_LOCK set
if (((s = runState) & RS_LOCK) == 0L && casRunState(s, u = s + RS_LOCK))
@@ -2015,15 +2027,19 @@ public class ForkJoinPool extends AbstractExecutorService
*/
final void runWorker(WorkQueue w) {
if (w != null) {
- int phase = w.phase, r = w.stackPred, src = UNSCANNED;
+ int phase = w.phase, r = w.stackPred, src = UNSCANNED, rescans = 0;
for (long e; ((e = runState) & STOP) == 0L; ) {
long stat = scan(w, r, phase, src);
phase = (int)(stat & LMASK);
src = (int)(stat >>> 32);
r ^= r << 13; r ^= r >>> 17; r ^= r << 5; // xorshift
- if (src == EMPTY_SCAN) {
- if ((e & RS_LOCK) == 0L && runState == e) {
- if ((phase & IDLE) == 0) { // try to deactivate
+ if (src >= 0)
+ rescans = 2;
+ else {
+ if (src == EMPTY_SCAN && (e & RS_LOCK) == 0L && runState == e) {
+ if (rescans > 0)
+ --rescans;
+ else if ((phase & IDLE) == 0) { // try to deactivate
int ip = phase | IDLE, ap = phase + (IDLE << 1);
long pc = ctl;
long qc = ((pc - RC_UNIT) & UMASK) | (ap & LMASK);
@@ -2035,8 +2051,8 @@ public class ForkJoinPool extends AbstractExecutorService
else if (awaitWork(w, phase))
break;
}
- src = UNSCANNED;
phase = w.phase;
+ src = UNSCANNED;
}
}
}
@@ -2055,9 +2071,8 @@ public class ForkJoinPool extends AbstractExecutorService
if ((qs = queues) == null || w == null || (n = qs.length) <= 0)
src = EMPTY_SCAN;
else {
- int polls = (src == UNSCANNED) ? n : Math.max(n << 1, MAX_REPOLLS);
- int idle = phase & IDLE;
- outer: for (int i = r, stride = (r >>> 16) | 1; ; i += stride) {
+ int polls = n, i = r, stride = (r >>> 16) | 1, idle = phase & IDLE;
+ outer: for (; ; i += stride) {
WorkQueue q; int qid;
if ((q = qs[qid = i & (n - 1)]) != null) {
boolean taken = false;
@@ -2070,16 +2085,13 @@ public class ForkJoinPool extends AbstractExecutorService
(tj = m & b + 2) < 0 || tj >= cap)
break; // never true
ForkJoinTask> t = getSlotAcquire(a, bj);
- boolean noNext = (a[nj] == null && a[tj] == null);
- ForkJoinTask> rt = a[bj];
- if (q.array != a) // resized
- ;
- else if (q.base != b || rt != t) { // inconsistent / busy
+ ForkJoinTask> nt = a[nj], tt = a[tj], rt = a[bj];
+ if (q.base != b || rt != t) { // inconsistent / busy
if (src != qid)
break outer; // reduce interference
}
else if (t == null) {
- if (noNext) {
+ if (nt == null && tt == null) {
if (taken) // end run
break outer;
break; // probably empty
@@ -2101,23 +2113,9 @@ public class ForkJoinPool extends AbstractExecutorService
break outer; // lost CAS; reorder scan
}
else if ((idle = (phase = w.phase) & IDLE) != 0) {
- WorkQueue v; int sp, sj; long c;
src = UNSCANNED; // possibly reactivate
- if ((r & 0xf00) == 0 &&
- (sp = (int)(c = ctl)) != 0 &&
- (v = qs[sp & (n - 1)]) != null) {
- long nc = ((v.stackPred & LMASK) |
- ((c + RC_UNIT) & UMASK));
- if (((phase = w.phase) & IDLE) != 0 &&
- (sj = m & q.base) >= 0 && sj < a.length &&
- a[sj] != null &&
- U.compareAndSetLong(this, CTL, c, nc)) {
- v.phase = sp;
- if (v != w && v.parking == sp)
- U.unpark(v.owner);
- phase = w.phase;
- }
- }
+ if ((r & 0xf00) == 0)
+ phase = tryReactivate(w, phase, qid);
break outer; // restart scan
}
}
@@ -2131,6 +2129,28 @@ public class ForkJoinPool extends AbstractExecutorService
return (((long)src) << 32) | (phase & LMASK);
}
+ private int tryReactivate(WorkQueue w, int phase, int qid) {
+ WorkQueue[] qs; WorkQueue v; int n, sp; long c;
+ if (w != null && (sp = (int)(c = ctl)) != 0 &&
+ (qs = queues) != null && (n = qs.length) > 0 &&
+ (v = qs[sp & (n - 1)]) != null) {
+ WorkQueue q; int cap, sj; ForkJoinTask>[] a;
+ long nc = ((v.stackPred & LMASK) | ((c + RC_UNIT) & UMASK));
+ if (qid >= 0 && qid < n && (q = qs[qid]) != null &&
+ (a = q.array) != null && (cap = a.length) > 0 &&
+ ((phase = w.phase) & IDLE) != 0 &&
+ (sj = (cap - 1) & q.base) >= 0 && sj < cap &&
+ getSlotAcquire(a, sj) != null &&
+ U.compareAndSetLong(this, CTL, c, nc)) {
+ v.phase = sp;
+ if (v != w && v.parking == sp)
+ U.unpark(v.owner);
+ phase = w.phase;
+ }
+ }
+ return phase;
+ }
+
/**
* Awaits signal or termination or spurious unpark.
*
@@ -2372,36 +2392,38 @@ public class ForkJoinPool extends AbstractExecutorService
int s = 0;
if (task != null && (s = task.status) >= 0 && internal && w != null) {
int wid = w.phase & SMASK, wsrc = w.source, r = wid + 2;
- long sctl = 0L; // track stability
+ long sctl = 0L, e = 0L; // track stability
outer: for (boolean rescan = true, found = true;;) {
- WorkQueue[] qs; int n;
+ WorkQueue[] qs; int n; long rs;
if ((s = task.status) < 0)
break;
- if ((runState & STOP) != 0L)
+ if (((rs = runState) & STOP) != 0L)
break;
- if (!rescan && sctl == (sctl = ctl) && (s = tryCompensate(sctl)) >= 0)
+ if (e == (e = rs) && !rescan && sctl == (sctl = ctl) &&
+ (s = tryCompensate(sctl)) >= 0)
break;
if (!found)
r = ThreadLocalRandom.nextSecondarySeed() & ~1;
+ int stride = ((r >>> 16) & ~1) | 2;
rescan = found = false;
if ((qs = queues) != null && (n = qs.length) > 0) {
- scan: for (int polls = n << 2; ; r += 2) {
- int j; WorkQueue q;
- if ((q = qs[j = r & SMASK & (n - 1)]) != null) {
- for (int b = q.base, pb = b - 1; ; b = q.base) {
- ForkJoinTask>[] a;
- boolean eligible = false;
- int sq = q.source, cap, m, nb, bj, nj;
- if ((a = q.array) == null || (cap = a.length) <= 0)
- break;
- if ((bj = (m = cap - 1) & b) < 0 || bj >= cap ||
+ scan: for (int polls = n << 1; ; r += stride) {
+ int j; WorkQueue q; ForkJoinTask>[] a; int cap;
+ if ((q = qs[j = r & SMASK & (n - 1)]) != null &&
+ (a = q.array) != null && (cap = a.length) > 0) {
+ int m = cap - 1, b = q.base, pb = b - 1;
+ for (; ; b = q.base) {
+ int nb, bj, nj;
+ if ((bj = m & b) < 0 || bj >= cap ||
(nj = m & (nb = b + 1)) < 0 || nj >= cap)
break; // never true
+ int sq = q.source;
+ boolean eligible = false;
ForkJoinTask> t = getSlotAcquire(a, bj);
if (t == task)
eligible = true;
else if (t != null) { // check steal chain
- for (int v = sq, d = cap;;) {
+ for (int v = sq, d = INITIAL_QUEUE_CAPACITY;;) {
WorkQueue p;
if (v == wid) {
eligible = true;
@@ -2433,7 +2455,7 @@ public class ForkJoinPool extends AbstractExecutorService
if (casSlotNull(a, bj, t)) {
q.base = nb;
w.source = j;
- t.doExec();
+ doExec(t);;
w.source = wsrc;
rescan = found = true; // restart at index r
break scan;
@@ -2464,33 +2486,36 @@ public class ForkJoinPool extends AbstractExecutorService
int s = 0;
if (task != null && (s = task.status) >= 0 && w != null) {
int r = w.phase + 1; // for indexing
- long sctl = 0L; // track stability
+ long sctl = 0L, e = 0L; // track stability
outer: for (boolean rescan = true, locals = true;;) {
- WorkQueue[] qs; int n;
- if (locals && (s = w.helpComplete(task, internal, 0)) < 0)
- break;
+ WorkQueue[] qs; int n; long rs;
if ((s = task.status) < 0)
break;
- if ((runState & STOP) != 0L)
+ if (((rs = runState) & STOP) != 0L)
break;
- if (!rescan && sctl == (sctl = ctl) &&
+ if (locals &&
+ ((s = w.helpComplete(task, internal, 0)) < 0 ||
+ (s = task.status) < 0))
+ break;
+ if (e == (e = rs) && !rescan && sctl == (sctl = ctl) &&
(!internal || (s = tryCompensate(sctl)) >= 0))
break;
if (!locals)
r = ThreadLocalRandom.nextSecondarySeed();
+ int stride = (r >>> 16) | 1;
rescan = locals = false;
if ((qs = queues) != null && (n = qs.length) > 0) {
- scan: for (int polls = n << 2; ; ++r) {
- int j; WorkQueue q;
- if ((q = qs[j = r & SMASK & (n - 1)]) != null) {
- for (int b = q.base, pb = b - 1; ; b = q.base) {
- ForkJoinTask>[] a; int cap, m, nb, bj, nj;
- boolean eligible = false;
- if ((a = q.array) == null || (cap = a.length) <= 0)
- break;
- if ((bj = (m = cap - 1) & b) < 0 || bj >= cap ||
+ scan: for (int polls = n << 1; ; r += stride) {
+ int j; WorkQueue q; ForkJoinTask>[] a; int cap;
+ if ((q = qs[j = r & SMASK & (n - 1)]) != null &&
+ (a = q.array) != null && (cap = a.length) > 0) {
+ int m = cap - 1, b = q.base, pb = b - 1;
+ for (; ; b = q.base) {
+ int nb, bj, nj;
+ if ((bj = m & b) < 0 || bj >= cap ||
(nj = m & (nb = b + 1)) < 0 || nj >= cap)
break; // never true
+ boolean eligible = false;
ForkJoinTask> t = getSlotAcquire(a, bj);
if (t != null && t instanceof CountedCompleter) {
CountedCompleter> f = (CountedCompleter>)t;
@@ -2522,7 +2547,7 @@ public class ForkJoinPool extends AbstractExecutorService
if (casSlotNull(a, bj, t)) {
q.base = nb;
U.storeFence();
- t.doExec();
+ doExec(t);;
locals = rescan = true;
break scan;
}
@@ -2572,19 +2597,20 @@ public class ForkJoinPool extends AbstractExecutorService
r = ThreadLocalRandom.nextSecondarySeed();
else {
for (ForkJoinTask> u; (u = w.nextLocalTask()) != null;)
- u.doExec(); // run local tasks before (re)polling
+ doExec(u); // run local tasks before (re)polling
}
long phaseSum = 0L;
+ int stride = (r >>> 16) | 1;
boolean rescan = false, busy = locals = false;
if ((qs = queues) != null && (n = qs.length) > 0) {
- scan: for (int polls = n; ; ++r) {
- int j; WorkQueue q;
- if ((q = qs[j = r & SMASK & (n - 1)]) != null && q != w) {
- for (int b = q.base, pb = b - 1; ; b = q.base) {
- ForkJoinTask>[] a; int cap, m, nb, bj, nj;
- if ((a = q.array) == null || (cap = a.length) <= 0)
- break;
- if ((bj = (m = cap - 1) & b) < 0 || bj >= cap ||
+ scan: for (int polls = n; ; r += stride) {
+ int j; WorkQueue q; ForkJoinTask>[] a; int cap;
+ if ((q = qs[j = r & SMASK & (n - 1)]) != null &&
+ (a = q.array) != null && (cap = a.length) > 0) {
+ int m = cap - 1, b = q.base, pb = b - 1;
+ for (; ; b = q.base) {
+ int nb, bj, nj;
+ if ((bj = m & b) < 0 || bj >= cap ||
(nj = m & (nb = b + 1)) < 0 || nj >= cap)
break; // never true
ForkJoinTask> t = getSlotAcquire(a, bj);
@@ -2613,7 +2639,7 @@ public class ForkJoinPool extends AbstractExecutorService
if (casSlotNull(a, bj, t)) {
q.base = nb;
w.source = j;
- t.doExec();
+ doExec(t);;
w.source = wsrc;
rescan = locals = true;
break scan;
@@ -2669,7 +2695,7 @@ public class ForkJoinPool extends AbstractExecutorService
return -1;
else if ((t = pollScan(false)) != null) {
waits = 0;
- t.doExec();
+ doExec(t);;
}
else if (quiescent() >= 0)
break;