Merge branch '6.2.x'
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2024 the original author or authors.
|
||||
* Copyright 2002-2025 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -87,6 +87,8 @@ public class SimpleAsyncTaskExecutor extends CustomizableThreadCreator
|
||||
|
||||
private @Nullable Set<Thread> activeThreads;
|
||||
|
||||
private boolean rejectTasksWhenLimitReached = false;
|
||||
|
||||
private volatile boolean active = true;
|
||||
|
||||
|
||||
@@ -184,6 +186,17 @@ public class SimpleAsyncTaskExecutor extends CustomizableThreadCreator
|
||||
this.activeThreads = (timeout > 0 ? ConcurrentHashMap.newKeySet() : null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify whether to reject tasks when the concurrency limit has been reached,
|
||||
* throwing {@link TaskRejectedException} on any further submission attempts.
|
||||
* <p>The default is {@code false}, blocking the caller until the submission can
|
||||
* be accepted. Switch this to {@code true} for immediate rejection instead.
|
||||
* @since 6.2.6
|
||||
*/
|
||||
public void setRejectTasksWhenLimitReached(boolean rejectTasksWhenLimitReached) {
|
||||
this.rejectTasksWhenLimitReached = rejectTasksWhenLimitReached;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set the maximum number of parallel task executions allowed.
|
||||
* The default of -1 indicates no concurrency limit at all.
|
||||
@@ -350,13 +363,21 @@ public class SimpleAsyncTaskExecutor extends CustomizableThreadCreator
|
||||
* making {@code beforeAccess()} and {@code afterAccess()}
|
||||
* visible to the surrounding class.
|
||||
*/
|
||||
private static class ConcurrencyThrottleAdapter extends ConcurrencyThrottleSupport {
|
||||
private class ConcurrencyThrottleAdapter extends ConcurrencyThrottleSupport {
|
||||
|
||||
@Override
|
||||
protected void beforeAccess() {
|
||||
super.beforeAccess();
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void onLimitReached() {
|
||||
if (rejectTasksWhenLimitReached) {
|
||||
throw new TaskRejectedException("Concurrency limit reached: " + getConcurrencyLimit());
|
||||
}
|
||||
super.onLimitReached();
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void afterAccess() {
|
||||
super.afterAccess();
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2024 the original author or authors.
|
||||
* Copyright 2002-2025 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -105,6 +105,7 @@ public abstract class ConcurrencyThrottleSupport implements Serializable {
|
||||
/**
|
||||
* To be invoked before the main execution logic of concrete subclasses.
|
||||
* <p>This implementation applies the concurrency throttle.
|
||||
* @see #onLimitReached()
|
||||
* @see #afterAccess()
|
||||
*/
|
||||
protected void beforeAccess() {
|
||||
@@ -113,29 +114,12 @@ public abstract class ConcurrencyThrottleSupport implements Serializable {
|
||||
"Currently no invocations allowed - concurrency limit set to NO_CONCURRENCY");
|
||||
}
|
||||
if (this.concurrencyLimit > 0) {
|
||||
boolean debug = logger.isDebugEnabled();
|
||||
this.concurrencyLock.lock();
|
||||
try {
|
||||
boolean interrupted = false;
|
||||
while (this.concurrencyCount >= this.concurrencyLimit) {
|
||||
if (interrupted) {
|
||||
throw new IllegalStateException("Thread was interrupted while waiting for invocation access, " +
|
||||
"but concurrency limit still does not allow for entering");
|
||||
}
|
||||
if (debug) {
|
||||
logger.debug("Concurrency count " + this.concurrencyCount +
|
||||
" has reached limit " + this.concurrencyLimit + " - blocking");
|
||||
}
|
||||
try {
|
||||
this.concurrencyCondition.await();
|
||||
}
|
||||
catch (InterruptedException ex) {
|
||||
// Re-interrupt current thread, to allow other threads to react.
|
||||
Thread.currentThread().interrupt();
|
||||
interrupted = true;
|
||||
}
|
||||
if (this.concurrencyCount >= this.concurrencyLimit) {
|
||||
onLimitReached();
|
||||
}
|
||||
if (debug) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Entering throttle at concurrency count " + this.concurrencyCount);
|
||||
}
|
||||
this.concurrencyCount++;
|
||||
@@ -146,6 +130,33 @@ public abstract class ConcurrencyThrottleSupport implements Serializable {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Triggered by {@link #beforeAccess()} when the concurrency limit has been reached.
|
||||
* The default implementation blocks until the concurrency count allows for entering.
|
||||
* @since 6.2.6
|
||||
*/
|
||||
protected void onLimitReached() {
|
||||
boolean interrupted = false;
|
||||
while (this.concurrencyCount >= this.concurrencyLimit) {
|
||||
if (interrupted) {
|
||||
throw new IllegalStateException("Thread was interrupted while waiting for invocation access, " +
|
||||
"but concurrency limit still does not allow for entering");
|
||||
}
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Concurrency count " + this.concurrencyCount +
|
||||
" has reached limit " + this.concurrencyLimit + " - blocking");
|
||||
}
|
||||
try {
|
||||
this.concurrencyCondition.await();
|
||||
}
|
||||
catch (InterruptedException ex) {
|
||||
// Re-interrupt current thread, to allow other threads to react.
|
||||
Thread.currentThread().interrupt();
|
||||
interrupted = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* To be invoked after the main execution logic of concrete subclasses.
|
||||
* @see #beforeAccess()
|
||||
|
||||
Reference in New Issue
Block a user