Add task rejection support to SyncTaskExecutor's concurrency throttle

Closes gh-36114
This commit is contained in:
Juergen Hoeller
2026-01-08 15:51:17 +01:00
parent 38f5f4de8e
commit c4d5e3c57d
3 changed files with 65 additions and 15 deletions
@@ -25,6 +25,7 @@ import java.util.concurrent.atomic.AtomicInteger;
import org.junit.jupiter.api.Test;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
import static org.assertj.core.api.Assertions.assertThatIOException;
import static org.assertj.core.api.Assertions.assertThatNoException;
@@ -36,23 +37,23 @@ class SyncTaskExecutorTests {
@Test
void plainExecution() {
SyncTaskExecutor taskExecutor = new SyncTaskExecutor();
SyncTaskExecutor executor = new SyncTaskExecutor();
ConcurrentClass target = new ConcurrentClass();
assertThatNoException().isThrownBy(() -> taskExecutor.execute(target::concurrentOperation));
assertThat(taskExecutor.execute(target::concurrentOperationWithResult)).isEqualTo("result");
assertThatIOException().isThrownBy(() -> taskExecutor.execute(target::concurrentOperationWithException));
assertThatNoException().isThrownBy(() -> executor.execute(target::concurrentOperation));
assertThat(executor.execute(target::concurrentOperationWithResult)).isEqualTo("result");
assertThatIOException().isThrownBy(() -> executor.execute(target::concurrentOperationWithException));
}
@Test
void withConcurrencyLimit() {
SyncTaskExecutor taskExecutor = new SyncTaskExecutor();
taskExecutor.setConcurrencyLimit(2);
SyncTaskExecutor executor = new SyncTaskExecutor();
executor.setConcurrencyLimit(2);
ConcurrentClass target = new ConcurrentClass();
List<CompletableFuture<?>> futures = new ArrayList<>(10);
for (int i = 0; i < 10; i++) {
futures.add(CompletableFuture.runAsync(() -> taskExecutor.execute(target::concurrentOperation)));
futures.add(CompletableFuture.runAsync(() -> executor.execute(target::concurrentOperation)));
}
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
assertThat(target.current).hasValue(0);
@@ -61,14 +62,14 @@ class SyncTaskExecutorTests {
@Test
void withConcurrencyLimitAndResult() {
SyncTaskExecutor taskExecutor = new SyncTaskExecutor();
taskExecutor.setConcurrencyLimit(2);
SyncTaskExecutor executor = new SyncTaskExecutor();
executor.setConcurrencyLimit(2);
ConcurrentClass target = new ConcurrentClass();
List<CompletableFuture<?>> futures = new ArrayList<>(10);
for (int i = 0; i < 10; i++) {
futures.add(CompletableFuture.runAsync(() ->
assertThat(taskExecutor.execute(target::concurrentOperationWithResult)).isEqualTo("result")));
assertThat(executor.execute(target::concurrentOperationWithResult)).isEqualTo("result")));
}
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
assertThat(target.current).hasValue(0);
@@ -77,20 +78,41 @@ class SyncTaskExecutorTests {
@Test
void withConcurrencyLimitAndException() {
SyncTaskExecutor taskExecutor = new SyncTaskExecutor();
taskExecutor.setConcurrencyLimit(2);
SyncTaskExecutor executor = new SyncTaskExecutor();
executor.setConcurrencyLimit(2);
ConcurrentClass target = new ConcurrentClass();
List<CompletableFuture<?>> futures = new ArrayList<>(10);
for (int i = 0; i < 10; i++) {
futures.add(CompletableFuture.runAsync(() ->
assertThatIOException().isThrownBy(() -> taskExecutor.execute(target::concurrentOperationWithException))));
assertThatIOException().isThrownBy(() -> executor.execute(target::concurrentOperationWithException))));
}
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
assertThat(target.current).hasValue(0);
assertThat(target.counter).hasValue(10);
}
@Test
void taskRejectedWhenConcurrencyLimitReached() throws Exception {
SyncTaskExecutor executor = new SyncTaskExecutor();
executor.setConcurrencyLimit(2);
executor.setRejectTasksWhenLimitReached(true);
ConcurrentClass target = new ConcurrentClass();
List<CompletableFuture<?>> futures = new ArrayList<>(10);
for (int i = 0; i < 2; i++) {
futures.add(CompletableFuture.runAsync(() -> executor.execute(target::concurrentOperation)));
}
Thread.sleep(10);
for (int i = 2; i < 10; i++) {
futures.add(CompletableFuture.runAsync(() ->
assertThatExceptionOfType(TaskRejectedException.class).isThrownBy(() -> executor.execute(target::concurrentOperation))));
}
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
assertThat(target.current).hasValue(0);
assertThat(target.counter).hasValue(2);
}
static class ConcurrentClass {
@@ -103,7 +125,7 @@ class SyncTaskExecutorTests {
throw new IllegalStateException();
}
try {
Thread.sleep(10);
Thread.sleep(100);
}
catch (InterruptedException ex) {
throw new IllegalStateException(ex);