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 5d24d1da70 fix(logging): isolate selector consume failures (#7128)
5d24d1da70 is described below

commit 5d24d1da70b3fbe634589551d5fa563474233b9c
Author: Liming Deng <[email protected]>
AuthorDate: Sun Sep 20 09:56:10 2026 +0800

    fix(logging): isolate selector consume failures (#7128)
    
    Co-authored-by: aias00 <[email protected]>
---
 .../common/collector/AbstractLogCollector.java     | 22 ++++----
 .../common/collector/AbstractLogCollectorTest.java | 60 ++++++++++++++++++++++
 2 files changed, 73 insertions(+), 9 deletions(-)

diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/main/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollector.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/main/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollector.java
index 560ff13c94..32116b0b60 100644
--- 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/main/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollector.java
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/main/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollector.java
@@ -114,15 +114,7 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
                 List<L> logs = new ArrayList<>();
                 int batchSize = 100;
                 if (getMultiClient()) {
-                    bufferQueueS.forEach((selectorId, bufferQueue) -> {
-                        List<L> logsS = new ArrayList<>();
-                        Long lastPushTime = lastPushTimeS.get(selectorId);
-                        try {
-                            processBufferQueue(bufferQueue, batchSize, 
diffTimeMSForPush, logsS, lastPushTime, selectorId);
-                        } catch (Exception e) {
-                            throw new RuntimeException(e);
-                        }
-                    });
+                    processMultiClientBufferQueues(batchSize, 
diffTimeMSForPush);
                 } else {
                     processBufferQueue(bufferQueue, batchSize, 
diffTimeMSForPush, logs, lastPushTime);
                 }
@@ -133,6 +125,18 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
         }
     }
 
+    void processMultiClientBufferQueues(final int batchSize, final int 
diffTimeMSForPush) {
+        bufferQueueS.forEach((selectorId, bufferQueue) -> {
+            List<L> logs = new ArrayList<>();
+            Long lastPushTime = lastPushTimeS.get(selectorId);
+            try {
+                processBufferQueue(bufferQueue, batchSize, diffTimeMSForPush, 
logs, lastPushTime, selectorId);
+            } catch (Exception e) {
+                LOG.error("Log collector failed to consume logs for selector 
{}", selectorId, e);
+            }
+        });
+    }
+
     private BlockingQueue<L> initQueue(final String selectorId) {
         bufferSize = getLogCollectConfig().getBufferQueueSize();
         bufferQueue = new LinkedBlockingDeque<>(bufferSize);
diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/test/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollectorTest.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/test/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollectorTest.java
index b6db3e0893..5883f8acaa 100644
--- 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/test/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollectorTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/test/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollectorTest.java
@@ -27,6 +27,7 @@ import org.junit.jupiter.api.Test;
 
 import java.lang.reflect.Field;
 import java.util.Collections;
+import java.util.HashMap;
 import java.util.HashSet;
 import java.util.Map;
 import java.util.Set;
@@ -38,7 +39,10 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotEquals;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.ArgumentMatchers.anyList;
+import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
 
 /**
  * The Test Case For AbstractLogCollector.
@@ -123,6 +127,55 @@ public class AbstractLogCollectorTest {
         assertSame(bufferedLog, bufferQueue.peek());
     }
 
+    @Test
+    public void testMultiClientConsumeContinuesAfterSelectorFailure() throws 
Exception {
+        AbstractLogConsumeClient<?, ShenyuRequestLog> failingClient = 
mock(AbstractLogConsumeClient.class);
+        AbstractLogConsumeClient<?, ShenyuRequestLog> succeedingClient = 
mock(AbstractLogConsumeClient.class);
+        Map<String, AbstractLogConsumeClient<?, ShenyuRequestLog>> clients = 
new HashMap<>();
+        clients.put("failing", failingClient);
+        clients.put("succeeding", succeedingClient);
+        AbstractLogCollector<AbstractLogConsumeClient<?, ShenyuRequestLog>, 
ShenyuRequestLog, GenericGlobalConfig> multiClientCollector =
+                new AbstractLogCollector<>() {
+                    @Override
+                    protected AbstractLogConsumeClient<?, ShenyuRequestLog> 
getLogConsumeClient() {
+                        return logConsumeClient;
+                    }
+
+                    @Override
+                    protected AbstractLogConsumeClient<?, ShenyuRequestLog> 
getLogConsumeClient(final String selectorId) {
+                        return clients.get(selectorId);
+                    }
+
+                    @Override
+                    protected boolean getMultiClient() {
+                        return true;
+                    }
+
+                    @Override
+                    protected GenericGlobalConfig getLogCollectConfig() {
+                        return null;
+                    }
+
+                    @Override
+                    protected void desensitizeLog(final ShenyuRequestLog log, 
final KeyWordMatch keyWordMatch, final String desensitizeAlg) {
+                    }
+                };
+        BlockingQueue<ShenyuRequestLog> failingQueue = new 
LinkedBlockingDeque<>(1);
+        BlockingQueue<ShenyuRequestLog> succeedingQueue = new 
LinkedBlockingDeque<>(1);
+        failingQueue.add(new ShenyuRequestLog());
+        succeedingQueue.add(new ShenyuRequestLog());
+        getBufferQueues(multiClientCollector).put("failing", failingQueue);
+        getBufferQueues(multiClientCollector).put("succeeding", 
succeedingQueue);
+        getLastPushTimes(multiClientCollector).put("failing", 
System.currentTimeMillis());
+        getLastPushTimes(multiClientCollector).put("succeeding", 
System.currentTimeMillis());
+        doThrow(new Exception("backend 
unavailable")).when(failingClient).consume(anyList());
+
+        multiClientCollector.processMultiClientBufferQueues(1, 100);
+
+        verify(failingClient).consume(anyList());
+        verify(succeedingClient).consume(anyList());
+    }
+
     @Test
     public void testDesensitizeToleratesNullBoxedNumericFields() {
         // a chunked byte-type response reaches desensitize with 
responseContentLength,
@@ -165,6 +218,13 @@ public class AbstractLogCollectorTest {
         return (Map<String, BlockingQueue<ShenyuRequestLog>>) 
field.get(target);
     }
 
+    @SuppressWarnings("unchecked")
+    private static Map<String, Long> getLastPushTimes(final Object target) 
throws Exception {
+        Field field = 
AbstractLogCollector.class.getDeclaredField("lastPushTimeS");
+        field.setAccessible(true);
+        return (Map<String, Long>) field.get(target);
+    }
+
     private static final class StaleSizeLinkedBlockingDeque extends 
LinkedBlockingDeque<ShenyuRequestLog> {
 
         private StaleSizeLinkedBlockingDeque() {

Reply via email to