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 891bc5deb9 fix(logging): make collector startup idempotent (#7227)
891bc5deb9 is described below
commit 891bc5deb969591224e2a6416f30ecbc39518116
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 12:18:14 2026 +0800
fix(logging): make collector startup idempotent (#7227)
---
.../common/collector/AbstractLogCollector.java | 52 ++++++++++---
.../common/collector/AbstractLogCollectorTest.java | 88 +++++++++++++++++++++-
2 files changed, 128 insertions(+), 12 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 535022155a..3576996a76 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
@@ -57,9 +57,9 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
private static final Logger LOG =
LoggerFactory.getLogger(AbstractLogCollector.class);
- private int bufferSize;
+ private volatile int bufferSize;
- private BlockingQueue<L> bufferQueue;
+ private volatile BlockingQueue<L> bufferQueue;
private final Map<String, BlockingQueue<L>> bufferQueueS =
Maps.newConcurrentMap();
@@ -67,12 +67,18 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
private long lastPushTime;
- private final AtomicBoolean started = new AtomicBoolean(true);
+ private final AtomicBoolean started = new AtomicBoolean(false);
+
+ private ShenyuThreadPoolExecutor executor;
@Override
- public void start() {
- bufferSize = getLogCollectConfig().getBufferQueueSize();
- bufferQueue = new LinkedBlockingDeque<>(bufferSize);
+ public synchronized void start() {
+ refreshBufferSize();
+ if (started.get()) {
+ return;
+ }
+ // Admission is bounded in collect(), allowing live capacity changes
without dropping queued logs.
+ bufferQueue = new LinkedBlockingDeque<>();
ShenyuConfig config =
Optional.ofNullable(Singleton.INST.get(ShenyuConfig.class)).orElse(new
ShenyuConfig());
final ShenyuConfig.SharedPool sharedPool = config.getSharedPool();
ShenyuThreadPoolExecutor threadExecutor = new
ShenyuThreadPoolExecutor(sharedPool.getCorePoolSize(),
@@ -81,7 +87,22 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
ShenyuThreadFactory.create(config.getSharedPool().getPrefix(),
true),
new ThreadPoolExecutor.AbortPolicy());
started.set(true);
- threadExecutor.execute(this::consume);
+ executor = threadExecutor;
+ try {
+ threadExecutor.execute(() -> consume(threadExecutor));
+ } catch (RuntimeException e) {
+ started.set(false);
+ threadExecutor.shutdownNow();
+ throw e;
+ }
+ }
+
+ private void refreshBufferSize() {
+ int configuredSize = getLogCollectConfig().getBufferQueueSize();
+ if (configuredSize <= 0) {
+ throw new IllegalArgumentException("bufferQueueSize must be
positive");
+ }
+ bufferSize = configuredSize;
}
@Override
@@ -94,7 +115,13 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
BlockingQueue<L> bufferQueue =
bufferQueueS.computeIfAbsent(selectorId, bufferQueueS -> initQueue(selectorId));
bufferQueue.offer(log);
} else {
- bufferQueue.offer(log);
+ BlockingQueue<L> queue = bufferQueue;
+ synchronized (queue) {
+ // On shrink, retain the backlog but reject new logs until it
falls below the new limit.
+ if (queue.size() < bufferSize) {
+ queue.offer(log);
+ }
+ }
}
}
@@ -107,8 +134,8 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
/**
* batch and async consume.
*/
- private void consume() {
- while (started.get()) {
+ private void consume(final ThreadPoolExecutor consumerExecutor) {
+ while (!consumerExecutor.isShutdown()) {
int diffTimeMSForPush = 100;
try {
List<L> logs = new ArrayList<>();
@@ -284,8 +311,11 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
protected abstract void desensitizeLog(L log, KeyWordMatch keyWordMatch,
String desensitizeAlg);
@Override
- public void close() throws Exception {
+ public synchronized void close() throws Exception {
started.set(false);
+ if (Objects.nonNull(executor)) {
+ executor.shutdownNow();
+ }
AbstractLogConsumeClient<?, ?> logCollectClient =
getLogConsumeClient();
try {
flushBufferQueues();
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 55540e804e..205682125f 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
@@ -32,18 +32,26 @@ import java.util.HashMap;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
+import java.util.List;
+import java.util.ArrayList;
+import java.util.stream.IntStream;
+import org.apache.shenyu.common.concurrent.ShenyuThreadPoolExecutor;
+import org.mockito.MockedConstruction;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
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.assertNotSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.Mockito.inOrder;
@@ -51,6 +59,9 @@ import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.times;
+import static org.mockito.ArgumentMatchers.any;
/**
* The Test Case For AbstractLogCollector.
@@ -59,6 +70,8 @@ public class AbstractLogCollectorTest {
private final AbstractLogConsumeClient<?, ShenyuRequestLog>
logConsumeClient = mock(AbstractLogConsumeClient.class);
+ private final GenericGlobalConfig collectorConfig = new
GenericGlobalConfig();
+
private final AbstractLogCollector<AbstractLogConsumeClient<?,
ShenyuRequestLog>, ShenyuRequestLog, GenericGlobalConfig> collector =
new AbstractLogCollector<>() {
@Override
@@ -68,7 +81,7 @@ public class AbstractLogCollectorTest {
@Override
protected GenericGlobalConfig getLogCollectConfig() {
- return null;
+ return collectorConfig;
}
@Override
@@ -76,6 +89,76 @@ public class AbstractLogCollectorTest {
}
};
+ @Test
+ public void testRepeatedStartAndRestartDoNotLeakConsumers() throws
Exception {
+ try {
+ collector.start();
+ ThreadPoolExecutor first = (ThreadPoolExecutor)
ReflectionTestUtils.getField(collector, "executor");
+ Object queue = ReflectionTestUtils.getField(collector,
"bufferQueue");
+ collector.start();
+ assertSame(first, ReflectionTestUtils.getField(collector,
"executor"));
+ assertSame(queue, ReflectionTestUtils.getField(collector,
"bufferQueue"));
+ collector.close();
+ collector.start();
+ assertNotSame(first, ReflectionTestUtils.getField(collector,
"executor"));
+ assertTrue(first.awaitTermination(2, TimeUnit.SECONDS));
+ } finally {
+ collector.close();
+ }
+ ThreadPoolExecutor last = (ThreadPoolExecutor)
ReflectionTestUtils.getField(collector, "executor");
+ assertTrue(last.awaitTermination(2, TimeUnit.SECONDS));
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testLiveCapacityChangesPreserveBacklogAndConsumer() throws
Exception {
+ try (MockedConstruction<ShenyuThreadPoolExecutor> executors =
mockConstruction(ShenyuThreadPoolExecutor.class)) {
+ collectorConfig.setBufferQueueSize(1);
+ collector.start();
+ final BlockingQueue<ShenyuRequestLog> queue =
(BlockingQueue<ShenyuRequestLog>) ReflectionTestUtils.getField(collector,
"bufferQueue");
+ final ShenyuRequestLog first = new ShenyuRequestLog();
+ final ShenyuRequestLog second = new ShenyuRequestLog();
+ final ShenyuRequestLog third = new ShenyuRequestLog();
+ first.setRequestUri("/first");
+ second.setRequestUri("/second");
+ third.setRequestUri("/third");
+ collector.collect(first);
+ collector.collect(second);
+ assertEquals(List.of(first), new ArrayList<>(queue));
+ collectorConfig.setBufferQueueSize(3);
+ collector.start();
+ collector.collect(second);
+ collector.collect(third);
+ collector.collect(new ShenyuRequestLog());
+ assertEquals(List.of(first, second, third), new
ArrayList<>(queue));
+ collectorConfig.setBufferQueueSize(1);
+ collector.start();
+ collector.collect(new ShenyuRequestLog());
+ assertEquals(List.of(first, second, third), new
ArrayList<>(queue));
+ queue.drainTo(new ArrayList<>());
+ collector.collect(first);
+ collector.collect(second);
+ assertEquals(List.of(first), new ArrayList<>(queue));
+ assertSame(queue, ReflectionTestUtils.getField(collector,
"bufferQueue"));
+ assertEquals(1, executors.constructed().size());
+ verify(executors.constructed().get(0),
times(1)).execute(any(Runnable.class));
+ collector.close();
+ }
+ }
+
+ @Test
+ public void testConcurrentProducersRespectLiveCapacity() throws Exception {
+ try (MockedConstruction<ShenyuThreadPoolExecutor> executors =
mockConstruction(ShenyuThreadPoolExecutor.class)) {
+ collectorConfig.setBufferQueueSize(10);
+ collector.start();
+ IntStream.range(0, 1000).parallel().forEach(index ->
collector.collect(new ShenyuRequestLog()));
+ BlockingQueue<?> queue = (BlockingQueue<?>)
ReflectionTestUtils.getField(collector, "bufferQueue");
+ assertEquals(10, queue.size());
+ assertEquals(1, executors.constructed().size());
+ collector.close();
+ }
+ }
+
@Test
public void testSelectorQueueInitializationDoesNotReplaceGlobalQueue()
throws Exception {
final GenericGlobalConfig config = new GenericGlobalConfig();
@@ -115,6 +198,7 @@ public class AbstractLogCollectorTest {
@Test
public void testCollectAddsLogWhenBufferQueueHasCapacity() throws
Exception {
+ ((AtomicBoolean) ReflectionTestUtils.getField(collector,
"started")).set(true);
BlockingQueue<ShenyuRequestLog> bufferQueue = new
LinkedBlockingDeque<>(1);
setField(collector, "bufferSize", 1);
setField(collector, "bufferQueue", bufferQueue);
@@ -127,6 +211,7 @@ public class AbstractLogCollectorTest {
@Test
public void testCollectDoesNotThrowWhenBufferQueueIsFull() throws
Exception {
+ ((AtomicBoolean) ReflectionTestUtils.getField(collector,
"started")).set(true);
ShenyuRequestLog bufferedLog = new ShenyuRequestLog();
BlockingQueue<ShenyuRequestLog> bufferQueue = new
StaleSizeLinkedBlockingDeque();
bufferQueue.add(bufferedLog);
@@ -164,6 +249,7 @@ public class AbstractLogCollectorTest {
BlockingQueue<ShenyuRequestLog> bufferQueue = new
StaleSizeLinkedBlockingDeque();
bufferQueue.add(bufferedLog);
setField(multiClientCollector, "bufferSize", 1);
+ ((AtomicBoolean) ReflectionTestUtils.getField(multiClientCollector,
"started")).set(true);
getBufferQueues(multiClientCollector).put("selector", bufferQueue);
ShenyuRequestLog log = new ShenyuRequestLog();
log.setSelectorId("selector");