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();
+        }
+    }
+}
+

Reply via email to