mirror of
https://github.com/openjdk/jdk.git
synced 2026-08-03 06:35:31 +00:00
Test more guesses about FMR slowdowns
This commit is contained in:
parent
875ac928d3
commit
bd7ce2aa2e
@ -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;
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user