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 e3990f1041 fix(logging): flush buffered logs on close (#7127)
e3990f1041 is described below
commit e3990f1041f39fd23fb913c61ddc8c70f57c4de6
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 07:05:19 2026 +0800
fix(logging): flush buffered logs on close (#7127)
---
.../common/collector/AbstractLogCollector.java | 35 ++++++++++-
.../common/collector/AbstractLogCollectorTest.java | 72 ++++++++++++++++++++++
2 files changed, 104 insertions(+), 3 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 03d2dae9f2..535022155a 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
@@ -86,7 +86,7 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
@Override
public void collect(final L log) {
- if (Objects.isNull(log) ||
Objects.isNull(getLogConsumeClient(log.getSelectorId()))) {
+ if (!started.get() || Objects.isNull(log) ||
Objects.isNull(getLogConsumeClient(log.getSelectorId()))) {
return;
}
if (getMultiClient()) {
@@ -181,6 +181,31 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
}
}
+ private void flushBufferQueues() throws Exception {
+ if (getMultiClient()) {
+ for (Map.Entry<String, BlockingQueue<L>> entry :
bufferQueueS.entrySet()) {
+ flushBufferQueue(entry.getValue(),
getLogConsumeClient(entry.getKey()));
+ }
+ } else {
+ flushBufferQueue(bufferQueue, getLogConsumeClient());
+ }
+ }
+
+ private void flushBufferQueue(final BlockingQueue<L> queue, final
AbstractLogConsumeClient<?, L> logConsumeClient) throws Exception {
+ if (Objects.isNull(queue) || Objects.isNull(logConsumeClient)) {
+ return;
+ }
+ int batchSize = 100;
+ while (!queue.isEmpty()) {
+ List<L> logs = new ArrayList<>(batchSize);
+ queue.drainTo(logs, batchSize);
+ if (logs.isEmpty()) {
+ return;
+ }
+ logConsumeClient.consume(logs);
+ }
+ }
+
private void desensitizeShenyuRequestLog(final L logInfo, final
KeyWordMatch keyWordMatch, final String desensitizedAlg) {
logInfo.setClientIp(desensitizeForSingleWord(GenericLoggingConstant.CLIENT_IP,
logInfo.getClientIp(), keyWordMatch, desensitizedAlg));
logInfo.setTimeLocal(desensitizeForSingleWord(GenericLoggingConstant.TIME_LOCAL,
logInfo.getTimeLocal(), keyWordMatch, desensitizedAlg));
@@ -262,8 +287,12 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
public void close() throws Exception {
started.set(false);
AbstractLogConsumeClient<?, ?> logCollectClient =
getLogConsumeClient();
- if (Objects.nonNull(logCollectClient)) {
- logCollectClient.close();
+ try {
+ flushBufferQueues();
+ } finally {
+ if (Objects.nonNull(logCollectClient)) {
+ logCollectClient.close();
+ }
}
}
}
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 7b5973cfdd..55540e804e 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
@@ -44,6 +44,9 @@ 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.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.argThat;
+import static org.mockito.Mockito.inOrder;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
@@ -247,6 +250,75 @@ public class AbstractLogCollectorTest {
assertEquals(15L, log.getUpstreamResponseTime());
}
+ @Test
+ public void testCloseFlushesBufferedLogsBeforeClosingClient() throws
Exception {
+ BlockingQueue<ShenyuRequestLog> bufferQueue = new
LinkedBlockingDeque<>(2);
+ ShenyuRequestLog first = new ShenyuRequestLog();
+ ShenyuRequestLog second = new ShenyuRequestLog();
+ bufferQueue.add(first);
+ bufferQueue.add(second);
+ setField(collector, "bufferQueue", bufferQueue);
+
+ collector.close();
+
+ org.mockito.InOrder closeOrder = inOrder(logConsumeClient);
+ closeOrder.verify(logConsumeClient).consume(argThat(logs ->
logs.size() == 2
+ && logs.get(0) == first && logs.get(1) == second));
+ closeOrder.verify(logConsumeClient).close();
+ assertTrue(bufferQueue.isEmpty());
+ }
+
+ @Test
+ public void testCloseFlushesEveryMultiClientBuffer() throws Exception {
+ AbstractLogConsumeClient<?, ShenyuRequestLog> firstClient =
mock(AbstractLogConsumeClient.class);
+ AbstractLogConsumeClient<?, ShenyuRequestLog> secondClient =
mock(AbstractLogConsumeClient.class);
+ Map<String, AbstractLogConsumeClient<?, ShenyuRequestLog>> clients =
new HashMap<>();
+ clients.put("first", firstClient);
+ clients.put("second", secondClient);
+ 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) {
+ }
+ };
+ ShenyuRequestLog first = new ShenyuRequestLog();
+ ShenyuRequestLog second = new ShenyuRequestLog();
+ BlockingQueue<ShenyuRequestLog> firstQueue = new
LinkedBlockingDeque<>(1);
+ BlockingQueue<ShenyuRequestLog> secondQueue = new
LinkedBlockingDeque<>(1);
+ firstQueue.add(first);
+ secondQueue.add(second);
+ getBufferQueues(multiClientCollector).put("first", firstQueue);
+ getBufferQueues(multiClientCollector).put("second", secondQueue);
+
+ multiClientCollector.close();
+
+ verify(firstClient).consume(argThat(logs -> logs.size() == 1 &&
logs.get(0) == first));
+ verify(secondClient).consume(argThat(logs -> logs.size() == 1 &&
logs.get(0) == second));
+ verify(logConsumeClient).close();
+ assertTrue(firstQueue.isEmpty());
+ assertTrue(secondQueue.isEmpty());
+ }
+
private static void setField(final Object target, final String fieldName,
final Object value) throws Exception {
Field field = AbstractLogCollector.class.getDeclaredField(fieldName);
field.setAccessible(true);