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 17bf2f4cf9 fix(logging): backpressure saturated cloud callback
executors (#7275)
17bf2f4cf9 is described below
commit 17bf2f4cf9f9730bf8b8059b08330bc16d7221b7
Author: Liming Deng <[email protected]>
AuthorDate: Sun Sep 27 11:41:13 2026 +0800
fix(logging): backpressure saturated cloud callback executors (#7275)
Co-authored-by: aias00 <[email protected]>
---
.../sls/client/AliyunSlsLogCollectClient.java | 2 +-
.../client/CloudLogCallbackBackpressureTest.java | 69 ++++++++++++++++++++++
.../lts/client/HuaweiLtsLogCollectClient.java | 2 +-
.../client/CloudLogCallbackBackpressureTest.java | 66 +++++++++++++++++++++
.../cls/client/TencentClsLogCollectClient.java | 2 +-
.../client/CloudLogCallbackBackpressureTest.java | 66 +++++++++++++++++++++
6 files changed, 204 insertions(+), 3 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/main/java/org/apache/shenyu/plugin/aliyun/sls/client/AliyunSlsLogCollectClient.java
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/main/java/org/apache/shenyu/plugin/aliyun/sls/client/AliyunSlsLogCollectClient.java
index 0ca5b6da85..e66418e0ba 100644
---
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/main/java/org/apache/shenyu/plugin/aliyun/sls/client/AliyunSlsLogCollectClient.java
+++
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/main/java/org/apache/shenyu/plugin/aliyun/sls/client/AliyunSlsLogCollectClient.java
@@ -191,7 +191,7 @@ public class AliyunSlsLogCollectClient extends
AbstractLogConsumeClient<AliyunLo
}
return new ThreadPoolExecutor(sendThreadCount,
GenericLoggingConstant.MAX_ALLOW_THREADS, 60000L, TimeUnit.MILLISECONDS,
new
LinkedBlockingQueue<>(GenericLoggingConstant.MAX_QUEUE_NUMBER),
ShenyuThreadFactory.create("shenyu-aliyun-sls", true),
- new ThreadPoolExecutor.AbortPolicy());
+ new ThreadPoolExecutor.CallerRunsPolicy());
}
/**
diff --git
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/test/java/org/apache/shenyu/plugin/aliyun/sls/client/CloudLogCallbackBackpressureTest.java
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/test/java/org/apache/shenyu/plugin/aliyun/sls/client/CloudLogCallbackBackpressureTest.java
new file mode 100644
index 0000000000..cdaec6d880
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/test/java/org/apache/shenyu/plugin/aliyun/sls/client/CloudLogCallbackBackpressureTest.java
@@ -0,0 +1,69 @@
+/*
+ * 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.plugin.aliyun.sls.client;
+
+import org.apache.shenyu.plugin.aliyun.sls.config.AliyunLogCollectConfig;
+import org.junit.jupiter.api.Test;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class CloudLogCallbackBackpressureTest {
+
+ @Test
+ void saturatedPoolRunsCallbackOnSubmittingThread() throws Exception {
+ AliyunLogCollectConfig.AliyunSlsLogConfig config = new
AliyunLogCollectConfig.AliyunSlsLogConfig();
+ config.setSendThreadCount(1);
+ ThreadPoolExecutor executor =
ReflectionTestUtils.invokeMethod(AliyunSlsLogCollectClient.class,
"createThreadPoolExecutor", config);
+ assertEquals(60, executor.getKeepAliveTime(TimeUnit.SECONDS));
+ executor.setMaximumPoolSize(1);
+ CountDownLatch entered = new CountDownLatch(1);
+ CountDownLatch release = new CountDownLatch(1);
+ try {
+ executor.execute(() -> {
+ entered.countDown();
+ try {
+ release.await();
+ } catch (InterruptedException ex) {
+ Thread.currentThread().interrupt();
+ }
+ });
+ assertTrue(entered.await(5, TimeUnit.SECONDS));
+ int capacity = executor.getQueue().remainingCapacity();
+ for (int i = 0; i < capacity; i++) {
+ executor.execute(() -> { });
+ }
+ AtomicReference<Thread> callbackThread = new AtomicReference<>();
+ executor.execute(() -> callbackThread.set(Thread.currentThread()));
+ assertSame(Thread.currentThread(), callbackThread.get());
+ assertEquals(capacity, executor.getQueue().size());
+ } finally {
+ executor.shutdownNow();
+ release.countDown();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ }
+ }
+}
diff --git
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/main/java/org/apache/shenyu/plugin/huawei/lts/client/HuaweiLtsLogCollectClient.java
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/main/java/org/apache/shenyu/plugin/huawei/lts/client/HuaweiLtsLogCollectClient.java
index b0df92cda4..c6e2d0bfc4 100644
---
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/main/java/org/apache/shenyu/plugin/huawei/lts/client/HuaweiLtsLogCollectClient.java
+++
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/main/java/org/apache/shenyu/plugin/huawei/lts/client/HuaweiLtsLogCollectClient.java
@@ -177,7 +177,7 @@ public class HuaweiLtsLogCollectClient extends
AbstractLogConsumeClient<HuaweiLo
}
return new ThreadPoolExecutor(threadCount,
GenericLoggingConstant.MAX_ALLOW_THREADS, 60000L, TimeUnit.MILLISECONDS,
new
LinkedBlockingQueue<>(GenericLoggingConstant.MAX_QUEUE_NUMBER),
ShenyuThreadFactory.create("shenyu-huawei-lts", true),
- new ThreadPoolExecutor.AbortPolicy());
+ new ThreadPoolExecutor.CallerRunsPolicy());
}
/**
diff --git
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/test/java/org/apache/shenyu/plugin/huawei/lts/client/CloudLogCallbackBackpressureTest.java
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/test/java/org/apache/shenyu/plugin/huawei/lts/client/CloudLogCallbackBackpressureTest.java
new file mode 100644
index 0000000000..3b97cef9d8
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/test/java/org/apache/shenyu/plugin/huawei/lts/client/CloudLogCallbackBackpressureTest.java
@@ -0,0 +1,66 @@
+/*
+ * 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.plugin.huawei.lts.client;
+
+import org.junit.jupiter.api.Test;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class CloudLogCallbackBackpressureTest {
+
+ @Test
+ void saturatedPoolRunsCallbackOnSubmittingThread() throws Exception {
+ ThreadPoolExecutor executor =
ReflectionTestUtils.invokeMethod(HuaweiLtsLogCollectClient.class,
"createThreadPoolExecutor", 1);
+ assertEquals(60, executor.getKeepAliveTime(TimeUnit.SECONDS));
+ executor.setMaximumPoolSize(1);
+ CountDownLatch entered = new CountDownLatch(1);
+ CountDownLatch release = new CountDownLatch(1);
+ try {
+ executor.execute(() -> {
+ entered.countDown();
+ try {
+ release.await();
+ } catch (InterruptedException ex) {
+ Thread.currentThread().interrupt();
+ }
+ });
+ assertTrue(entered.await(5, TimeUnit.SECONDS));
+ int capacity = executor.getQueue().remainingCapacity();
+ for (int i = 0; i < capacity; i++) {
+ executor.execute(() -> { });
+ }
+ AtomicReference<Thread> callbackThread = new AtomicReference<>();
+ executor.execute(() -> callbackThread.set(Thread.currentThread()));
+ assertSame(Thread.currentThread(), callbackThread.get());
+ assertEquals(capacity, executor.getQueue().size());
+ } finally {
+ executor.shutdownNow();
+ release.countDown();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ }
+ }
+}
diff --git
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/main/java/org/apache/shenyu/plugin/tencent/cls/client/TencentClsLogCollectClient.java
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/main/java/org/apache/shenyu/plugin/tencent/cls/client/TencentClsLogCollectClient.java
index 4b90205033..ae82a573e3 100644
---
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/main/java/org/apache/shenyu/plugin/tencent/cls/client/TencentClsLogCollectClient.java
+++
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/main/java/org/apache/shenyu/plugin/tencent/cls/client/TencentClsLogCollectClient.java
@@ -170,7 +170,7 @@ public class TencentClsLogCollectClient extends
AbstractLogConsumeClient<Tencent
}
return new ThreadPoolExecutor(threadCount,
GenericLoggingConstant.MAX_ALLOW_THREADS, 60000L, TimeUnit.MILLISECONDS,
new
LinkedBlockingQueue<>(GenericLoggingConstant.MAX_QUEUE_NUMBER),
ShenyuThreadFactory.create("shenyu-tencent-cls", true),
- new ThreadPoolExecutor.AbortPolicy());
+ new ThreadPoolExecutor.CallerRunsPolicy());
}
/**
diff --git
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/test/java/org/apache/shenyu/plugin/tencent/cls/client/CloudLogCallbackBackpressureTest.java
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/test/java/org/apache/shenyu/plugin/tencent/cls/client/CloudLogCallbackBackpressureTest.java
new file mode 100644
index 0000000000..50eeec6870
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/test/java/org/apache/shenyu/plugin/tencent/cls/client/CloudLogCallbackBackpressureTest.java
@@ -0,0 +1,66 @@
+/*
+ * 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.plugin.tencent.cls.client;
+
+import org.junit.jupiter.api.Test;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class CloudLogCallbackBackpressureTest {
+
+ @Test
+ void saturatedPoolRunsCallbackOnSubmittingThread() throws Exception {
+ ThreadPoolExecutor executor =
ReflectionTestUtils.invokeMethod(TencentClsLogCollectClient.class,
"createThreadPoolExecutor", 1);
+ assertEquals(60, executor.getKeepAliveTime(TimeUnit.SECONDS));
+ executor.setMaximumPoolSize(1);
+ CountDownLatch entered = new CountDownLatch(1);
+ CountDownLatch release = new CountDownLatch(1);
+ try {
+ executor.execute(() -> {
+ entered.countDown();
+ try {
+ release.await();
+ } catch (InterruptedException ex) {
+ Thread.currentThread().interrupt();
+ }
+ });
+ assertTrue(entered.await(5, TimeUnit.SECONDS));
+ int capacity = executor.getQueue().remainingCapacity();
+ for (int i = 0; i < capacity; i++) {
+ executor.execute(() -> { });
+ }
+ AtomicReference<Thread> callbackThread = new AtomicReference<>();
+ executor.execute(() -> callbackThread.set(Thread.currentThread()));
+ assertSame(Thread.currentThread(), callbackThread.get());
+ assertEquals(capacity, executor.getQueue().size());
+ } finally {
+ executor.shutdownNow();
+ release.countDown();
+ assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+ }
+ }
+}