This is an automated email from the ASF dual-hosted git repository.
Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git
The following commit(s) were added to refs/heads/master by this push:
new dec991e5fd fix(disruptor): bound consumer queues and preserve ordered
backpressure (#7233)
dec991e5fd is described below
commit dec991e5fd2f215415cde96293170d24999f0e06
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 11:53:37 2026 +0800
fix(disruptor): bound consumer queues and preserve ordered backpressure
(#7233)
---
.../shenyu/disruptor/DisruptorProviderManage.java | 6 +-
.../disruptor/thread/BlockWhenFullPolicy.java | 47 ++++++++
.../shenyu/disruptor/thread/OrderlyExecutor.java | 2 +-
.../shenyu/disruptor/thread/SingletonExecutor.java | 8 +-
.../disruptor/thread/BlockWhenFullPolicyTest.java | 126 +++++++++++++++++++++
5 files changed, 184 insertions(+), 5 deletions(-)
diff --git
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/DisruptorProviderManage.java
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/DisruptorProviderManage.java
index 6512d01d49..84d3490ff6 100644
---
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/DisruptorProviderManage.java
+++
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/DisruptorProviderManage.java
@@ -31,9 +31,9 @@ import
org.apache.shenyu.disruptor.event.OrderlyDisruptorEventFactory;
import org.apache.shenyu.disruptor.provider.DisruptorProvider;
import org.apache.shenyu.disruptor.thread.DisruptorThreadFactory;
import org.apache.shenyu.disruptor.thread.OrderlyExecutor;
+import org.apache.shenyu.disruptor.thread.BlockWhenFullPolicy;
import java.util.concurrent.LinkedBlockingQueue;
-import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
/**
@@ -112,8 +112,8 @@ public class DisruptorProviderManage<T> {
return;
}
OrderlyExecutor executor = new OrderlyExecutor(isOrderly,
consumerSize, consumerSize, 0, TimeUnit.MILLISECONDS,
- new LinkedBlockingQueue<>(),
- DisruptorThreadFactory.create("shenyu_disruptor_consumer_",
false), new ThreadPoolExecutor.AbortPolicy());
+ new LinkedBlockingQueue<>(size),
+ DisruptorThreadFactory.create("shenyu_disruptor_consumer_",
false), new BlockWhenFullPolicy());
int newConsumerSize = this.consumerSize;
EventFactory<DataEvent<T>> eventFactory;
if (isOrderly) {
diff --git
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/BlockWhenFullPolicy.java
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/BlockWhenFullPolicy.java
new file mode 100644
index 0000000000..69b3a5f3d6
--- /dev/null
+++
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/BlockWhenFullPolicy.java
@@ -0,0 +1,47 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.disruptor.thread;
+
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.RejectedExecutionHandler;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * Apply backpressure to the Disruptor consumer without running ordered tasks
out of order.
+ */
+public final class BlockWhenFullPolicy implements RejectedExecutionHandler {
+
+ @Override
+ public void rejectedExecution(final Runnable task, final
ThreadPoolExecutor executor) {
+ try {
+ while (!executor.isShutdown()) {
+ if (executor.getQueue().offer(task, 100,
TimeUnit.MILLISECONDS)) {
+ if (executor.isShutdown() && executor.remove(task)) {
+ throw new RejectedExecutionException("Executor shut
down during handoff");
+ }
+ return;
+ }
+ }
+ throw new RejectedExecutionException("Executor is shut down");
+ } catch (InterruptedException exception) {
+ Thread.currentThread().interrupt();
+ throw new RejectedExecutionException("Interrupted while waiting
for consumer capacity", exception);
+ }
+ }
+}
diff --git
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/OrderlyExecutor.java
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/OrderlyExecutor.java
index d40287c063..ccc307d419 100644
---
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/OrderlyExecutor.java
+++
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/OrderlyExecutor.java
@@ -69,7 +69,7 @@ public class OrderlyExecutor extends ThreadPoolExecutor {
private void orderlyThreadPool(final boolean isOrderly, final int
corePoolSize, final ThreadFactory threadFactory) {
if (isOrderly) {
IntStream.range(0, corePoolSize).forEach(index -> {
- SingletonExecutor singletonExecutor = new
SingletonExecutor(threadFactory);
+ SingletonExecutor singletonExecutor = new
SingletonExecutor(threadFactory, getQueue().remainingCapacity());
String hash = singletonExecutor.hashCode() + ":" + index;
byte[] bytes = threadSelector.sha(hash);
for (int i = 0; i < 4; i++) {
diff --git
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/SingletonExecutor.java
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/SingletonExecutor.java
index 29a56ca957..2255d590af 100644
---
a/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/SingletonExecutor.java
+++
b/shenyu-disruptor/src/main/java/org/apache/shenyu/disruptor/thread/SingletonExecutor.java
@@ -17,6 +17,8 @@
package org.apache.shenyu.disruptor.thread;
+import org.apache.shenyu.disruptor.DisruptorProviderManage;
+
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
@@ -29,8 +31,12 @@ import java.util.concurrent.TimeUnit;
public class SingletonExecutor extends ThreadPoolExecutor {
public SingletonExecutor(final ThreadFactory factory) {
+ this(factory, DisruptorProviderManage.DEFAULT_SIZE);
+ }
+
+ public SingletonExecutor(final ThreadFactory factory, final int
queueCapacity) {
super(1, 1, 0L,
TimeUnit.MILLISECONDS,
- new LinkedBlockingQueue<>(), factory);
+ new LinkedBlockingQueue<>(queueCapacity), factory, new
BlockWhenFullPolicy());
}
}
diff --git
a/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/thread/BlockWhenFullPolicyTest.java
b/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/thread/BlockWhenFullPolicyTest.java
new file mode 100644
index 0000000000..e9338986e8
--- /dev/null
+++
b/shenyu-disruptor/src/test/java/org/apache/shenyu/disruptor/thread/BlockWhenFullPolicyTest.java
@@ -0,0 +1,126 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.shenyu.disruptor.thread;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.Executors;
+import java.util.concurrent.FutureTask;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class BlockWhenFullPolicyTest {
+
+ @Test
+ void testBackpressurePreservesOrder() throws Exception {
+ SingletonExecutor executor = new
SingletonExecutor(Executors.defaultThreadFactory(), 1);
+ CountDownLatch started = new CountDownLatch(1);
+ CountDownLatch release = new CountDownLatch(1);
+ List<Integer> order = new CopyOnWriteArrayList<>();
+ FutureTask<Void> submit = new FutureTask<>(() -> {
+ executor.execute(() -> order.add(3));
+ return null;
+ });
+ Thread producer = new Thread(submit);
+ try {
+ executor.execute(() -> {
+ started.countDown();
+ await(release);
+ order.add(1);
+ });
+ assertTrue(started.await(5, TimeUnit.SECONDS));
+ executor.execute(() -> order.add(2));
+ producer.start();
+ assertThrows(TimeoutException.class, () -> submit.get(100,
TimeUnit.MILLISECONDS));
+ assertEquals(1, executor.getQueue().size());
+ release.countDown();
+ submit.get(5, TimeUnit.SECONDS);
+ executor.shutdown();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ assertEquals(Arrays.asList(1, 2, 3), order);
+ } finally {
+ release.countDown();
+ producer.interrupt();
+ producer.join(5000);
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ }
+ }
+
+ @Test
+ void testShutdownReleasesBlockedProducer() throws Exception {
+ SingletonExecutor executor = new
SingletonExecutor(Executors.defaultThreadFactory(), 1);
+ CountDownLatch started = new CountDownLatch(1);
+ CountDownLatch release = new CountDownLatch(1);
+ FutureTask<Void> submit = new FutureTask<>(() -> {
+ assertThrows(RejectedExecutionException.class, () ->
executor.execute(() -> { }));
+ return null;
+ });
+ Thread producer = new Thread(submit);
+ try {
+ executor.execute(() -> {
+ started.countDown();
+ await(release);
+ });
+ assertTrue(started.await(5, TimeUnit.SECONDS));
+ executor.execute(() -> { });
+ producer.start();
+ assertThrows(TimeoutException.class, () -> submit.get(100,
TimeUnit.MILLISECONDS));
+ executor.shutdown();
+ submit.get(5, TimeUnit.SECONDS);
+ } finally {
+ release.countDown();
+ producer.interrupt();
+ producer.join(5000);
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ }
+ }
+
+ @Test
+ void testInterruptedProducerRestoresInterruptFlag() {
+ SingletonExecutor executor = new
SingletonExecutor(Executors.defaultThreadFactory(), 1);
+ executor.getQueue().add(() -> { });
+ try {
+ Thread.currentThread().interrupt();
+ assertThrows(RejectedExecutionException.class, () -> new
BlockWhenFullPolicy().rejectedExecution(() -> { }, executor));
+ assertTrue(Thread.currentThread().isInterrupted());
+ } finally {
+ Thread.interrupted();
+ executor.shutdownNow();
+ }
+ }
+
+ private void await(final CountDownLatch latch) {
+ try {
+ latch.await();
+ } catch (InterruptedException exception) {
+ Thread.currentThread().interrupt();
+ }
+ }
+}
+