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;