forked from dragonwell-project/dragonwell8_jdk
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
[Wisp] Thread-based asynchronous IO implementation
Summary: Use ThreadPool as asynchronous IO delegate for asynchronous IO implementation. The feature can be used in Wisp flow and non-Wisp flow. Test Plan: test/com/alibaba/wisp/io/ThreadPoolAIOTest.java Reviewed-by: leiyu, zhengxiaolinX Issue: dragonwell-project/dragonwell8#213 thread local judge
- Loading branch information
joeyleeeeeee97
authored and
joeylee.lz
committed
Apr 25, 2021
1 parent
2ee14c2
commit 81237ca
Showing
19 changed files
with
632 additions
and
26 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
121 changes: 121 additions & 0 deletions
121
src/share/classes/com/alibaba/wisp/engine/WispAIOSupporter.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,121 @@ | ||
package com.alibaba.wisp.engine; | ||
|
||
import java.io.IOException; | ||
import java.util.concurrent.*; | ||
import java.util.concurrent.atomic.AtomicInteger; | ||
import java.util.concurrent.TimeUnit; | ||
|
||
public enum WispAIOSupporter { | ||
INSTANCE; | ||
|
||
private ExecutorService executor; | ||
|
||
private ThreadGroup threadgroup; | ||
|
||
WispAIOSupporter() { | ||
} | ||
|
||
enum Policy { | ||
ABORT { | ||
@Override | ||
void handle(ThreadPoolExecutor executor) { | ||
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy()); | ||
} | ||
}, | ||
CALLERRUNS { | ||
@Override | ||
void handle(ThreadPoolExecutor executor) { | ||
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); | ||
} | ||
}, | ||
DISCARDOLDEST { | ||
@Override | ||
void handle(ThreadPoolExecutor executor) { | ||
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.DiscardOldestPolicy()); | ||
} | ||
}, | ||
DISCARD { | ||
@Override | ||
void handle(ThreadPoolExecutor executor) { | ||
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.DiscardPolicy()); | ||
} | ||
}; | ||
|
||
abstract void handle(ThreadPoolExecutor executor); | ||
} | ||
|
||
void startDaemon(ThreadGroup g) { | ||
threadgroup = g; | ||
ThreadPoolExecutor workPool; | ||
if (WispAsyncIO.AIO_QUEUE_LIMIT == -1) { | ||
workPool = new ThreadPoolExecutor(WispAsyncIO.CORE_POOL_SIZE, | ||
WispAsyncIO.MAX_POOL_SIZE, Long.MAX_VALUE, | ||
TimeUnit.SECONDS, new LinkedBlockingDeque<>(), new AIOThreadPoolFactory()); | ||
} else { | ||
workPool = new ThreadPoolExecutor(WispAsyncIO.CORE_POOL_SIZE, | ||
WispAsyncIO.MAX_POOL_SIZE, Long.MAX_VALUE, | ||
TimeUnit.SECONDS, new ArrayBlockingQueue<>(WispAsyncIO.AIO_QUEUE_LIMIT), new AIOThreadPoolFactory()); | ||
} | ||
WispAsyncIO.REJECTED_POLICY.handle(workPool); | ||
this.executor = workPool; | ||
WispAsyncIO.wisp_aio_done = true; | ||
} | ||
|
||
public <T> T invokeIOTask(Callable<T> command) throws IOException { | ||
Future<T> future; | ||
try { | ||
future = submitIOTask(command); | ||
} catch (RejectedExecutionException e) { | ||
throw new IOException("busy", e); | ||
} | ||
T result; | ||
while (true) { | ||
try { | ||
if (WispAsyncIO.TIMEOUT == -1) { | ||
result = future.get(); | ||
} else { | ||
result = future.get(WispAsyncIO.TIMEOUT, TimeUnit.MILLISECONDS); | ||
} | ||
return result; | ||
} catch (TimeoutException e) { | ||
throw new IOException(e); | ||
} catch (ExecutionException e) { | ||
Throwable cause = e.getCause(); | ||
Class<?> causeClass = cause.getClass(); | ||
if (IOException.class.isAssignableFrom(causeClass)) { | ||
throw (IOException) cause; | ||
} else if (RuntimeException.class.isAssignableFrom(causeClass)) { | ||
throw (RuntimeException) cause; | ||
} else if (Error.class.isAssignableFrom(causeClass)) { | ||
throw (Error) cause; | ||
} else { | ||
throw new Error(e); | ||
} | ||
} catch (InterruptedException e) { | ||
Thread.interrupted(); | ||
} | ||
} | ||
} | ||
|
||
<T> Future<T> submitIOTask(Callable<T> command) { | ||
return executor.submit(command); | ||
} | ||
|
||
private class AIOThreadPoolFactory implements ThreadFactory { | ||
private final AtomicInteger threadNumber = new AtomicInteger(1); | ||
private final static String namePrefix = "AIO-worker-thread-"; | ||
AIOThreadPoolFactory() {} | ||
|
||
@Override | ||
public Thread newThread(Runnable r) { | ||
Thread t; | ||
if (threadgroup != null) { | ||
t = new Thread(threadgroup, r, namePrefix + threadNumber.getAndIncrement()); | ||
} else { | ||
t = new Thread(r, namePrefix + threadNumber.getAndIncrement()); | ||
} | ||
t.setDaemon(true); | ||
return t; | ||
} | ||
} | ||
} |
105 changes: 105 additions & 0 deletions
105
src/share/classes/com/alibaba/wisp/engine/WispAsyncIO.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,105 @@ | ||
package com.alibaba.wisp.engine; | ||
|
||
import java.io.IOException; | ||
import java.util.concurrent.Callable; | ||
import java.util.Properties; | ||
import java.security.AccessController; | ||
import sun.misc.SharedSecrets; | ||
import sun.misc.UnsafeAccess; | ||
import sun.misc.WispAsyncIOAccess; | ||
import sun.security.action.GetBooleanSecurityPropertyAction; | ||
import sun.security.action.GetIntegerAction; | ||
|
||
public class WispAsyncIO { | ||
// AIO | ||
static final int CORE_POOL_SIZE; | ||
static final int MAX_POOL_SIZE; | ||
static final int AIO_QUEUE_LIMIT; | ||
static final int TIMEOUT; | ||
static final WispAIOSupporter.Policy REJECTED_POLICY; | ||
static boolean wisp_aio_done = false; | ||
|
||
static { | ||
Properties p = java.security.AccessController.doPrivileged( | ||
new java.security.PrivilegedAction<Properties>() { | ||
public Properties run() { | ||
return System.getProperties(); | ||
} | ||
} | ||
); | ||
CORE_POOL_SIZE = parsePositiveIntegerParameter(p, "com.alibaba.aio.corePoolSize", | ||
Runtime.getRuntime().availableProcessors()); | ||
MAX_POOL_SIZE = parsePositiveIntegerParameter(p, "com.alibaba.aio.maxPoolSize", | ||
Runtime.getRuntime().availableProcessors()); | ||
AIO_QUEUE_LIMIT = parsePositiveIntegerParameter(p, "com.alibaba.aio.queueLimit", -1); | ||
TIMEOUT = parsePositiveIntegerParameter(p, "com.alibaba.aio.timeout", -1); | ||
REJECTED_POLICY = WispAIOSupporter.Policy.valueOf( | ||
p.getProperty("com.alibaba.aio.policy", WispAIOSupporter.Policy.ABORT.name())); | ||
} | ||
|
||
static void setWispAsyncIOAccess() { | ||
// initialize TenantAccess | ||
if (SharedSecrets.getWispAsyncIOAccess() == null) { | ||
SharedSecrets.setWispAsyncIOAccess(new WispAsyncIOAccess() { | ||
@Override | ||
public boolean usingAsyncIO() { | ||
return wisp_aio_done; | ||
} | ||
|
||
@Override | ||
public <T> T executeAsyncIO(Callable<T> command) throws IOException { | ||
return WispAIOSupporter.INSTANCE.invokeIOTask(command); | ||
} | ||
}); | ||
} | ||
} | ||
|
||
public static boolean useAsyncIO() { | ||
return SharedSecrets.getWispAsyncIOAccess() != null | ||
&& SharedSecrets.getWispAsyncIOAccess().usingAsyncIO() | ||
&& SharedSecrets.getWispEngineAccess() != null | ||
&& SharedSecrets.getWispEngineAccess().runningAsCoroutine(Thread.currentThread()); | ||
} | ||
|
||
/* | ||
* Initialize the WispAsyncIO class, called after System.initializeSystemClass by VM. | ||
**/ | ||
static void initializeAsyncIOClass() { | ||
try { | ||
Class.forName(WispAIOSupporter.class.getName()); | ||
} catch (Exception e) { | ||
throw new ExceptionInInitializerError(e); | ||
} | ||
} | ||
|
||
static void startAsyncIODaemonWithThreadGroup(ThreadGroup g) { | ||
WispAIOSupporter.INSTANCE.startDaemon(g); | ||
} | ||
|
||
static void startAsyncIODaemon() { | ||
WispAIOSupporter.INSTANCE.startDaemon(WispEngine.DAEMON_THREAD_GROUP); | ||
setWispAsyncIOAccess(); | ||
} | ||
|
||
private static boolean parseBooleanParameter(Properties p, String key, boolean defaultVal) { | ||
String value; | ||
if (p == null || (value = p.getProperty(key)) == null) { | ||
return defaultVal; | ||
} | ||
return Boolean.valueOf(value); | ||
} | ||
|
||
private static int parsePositiveIntegerParameter(Properties p, String key, int defaultVal) { | ||
String value; | ||
if (p == null || (value = p.getProperty(key)) == null) { | ||
return defaultVal; | ||
} | ||
int res = defaultVal; | ||
try { | ||
res = Integer.valueOf(value); | ||
} catch (NumberFormatException e) { | ||
return defaultVal; | ||
} | ||
return res <= 0 ? defaultVal : res; | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.