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 474e88cf95 [type:fix] kafka logging: flush producer once per batch and
handle flush failure (#7043)
474e88cf95 is described below
commit 474e88cf95d1ce9a2095a5b42d337ff732bb4c90
Author: renxuefeng <[email protected]>
AuthorDate: Wed Sep 30 16:52:24 2026 +0800
[type:fix] kafka logging: flush producer once per batch and handle flush
failure (#7043)
Move producer.flush() out of the per-record loop so the producer can batch.
Wrap the batch level flush so a flush failure does not escape consume0 and make
AbstractLogCollector discard the whole drained batch. Add unit tests covering
the once-per-batch flush and the flush failure path.
---
.../kafka/client/KafkaLogCollectClient.java | 6 +++-
.../kafka/kafka/KafkaLogCollectClientTest.java | 32 +++++++++++++++++++++-
2 files changed, 36 insertions(+), 2 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-kafka/src/main/java/org/apache/shenyu/plugin/logging/kafka/client/KafkaLogCollectClient.java
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-kafka/src/main/java/org/apache/shenyu/plugin/logging/kafka/client/KafkaLogCollectClient.java
index e919111dab..241e347611 100644
---
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-kafka/src/main/java/org/apache/shenyu/plugin/logging/kafka/client/KafkaLogCollectClient.java
+++
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-kafka/src/main/java/org/apache/shenyu/plugin/logging/kafka/client/KafkaLogCollectClient.java
@@ -131,11 +131,15 @@ public class KafkaLogCollectClient extends
AbstractLogConsumeClient<KafkaLogColl
LOG.error("kafka push logs error", exception);
}
});
- producer.flush();
} catch (Exception e) {
LOG.error("kafka push logs error", e);
}
});
+ try {
+ producer.flush();
+ } catch (Exception e) {
+ LOG.error("kafka flush logs error", e);
+ }
}
private ProducerRecord<String, String> toProducerRecord(final String
logTopic, final ShenyuRequestLog log) {
diff --git
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-kafka/src/test/java/org/apache/shenyu/plugin/logging/kafka/kafka/KafkaLogCollectClientTest.java
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-kafka/src/test/java/org/apache/shenyu/plugin/logging/kafka/kafka/KafkaLogCollectClientTest.java
index a307fc3e4b..2a85aa4d2b 100644
---
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-kafka/src/test/java/org/apache/shenyu/plugin/logging/kafka/kafka/KafkaLogCollectClientTest.java
+++
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-kafka/src/test/java/org/apache/shenyu/plugin/logging/kafka/kafka/KafkaLogCollectClientTest.java
@@ -18,6 +18,7 @@
package org.apache.shenyu.plugin.logging.kafka.kafka;
import org.apache.kafka.clients.producer.KafkaProducer;
+import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.errors.TimeoutException;
import org.apache.shenyu.common.dto.PluginData;
import org.apache.shenyu.common.utils.GsonUtils;
@@ -29,9 +30,14 @@ import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.mockito.MockedConstruction;
+import java.lang.reflect.Field;
+import java.util.Collections;
+
import static org.mockito.ArgumentMatchers.any;
-import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -57,6 +63,7 @@ public class KafkaLogCollectClientTest {
globalLogConfig.setCompressAlg("LZ4");
shenyuRequestLog.setClientIp("0.0.0.0");
shenyuRequestLog.setPath("org/apache/shenyu/plugin/logging");
+ shenyuRequestLog.setSelectorId("test-selector-id");
}
@Test
@@ -79,4 +86,27 @@ public class KafkaLogCollectClientTest {
verify(construction.constructed().get(0), never()).send(any());
}
}
+
+ @Test
+ public void testConsume0FlushesOncePerBatch() throws NoSuchFieldException,
IllegalAccessException {
+ KafkaProducer<String, String> producer = mock(KafkaProducer.class);
+ setProducer(producer);
+
kafkaLogCollectClient.consume0(Collections.singletonList(shenyuRequestLog));
+ verify(producer).flush();
+ }
+
+ @Test
+ public void testConsume0HandlesFlushFailure() throws NoSuchFieldException,
IllegalAccessException {
+ KafkaProducer<String, String> producer = mock(KafkaProducer.class);
+ doThrow(new KafkaException("flush error")).when(producer).flush();
+ setProducer(producer);
+ Assertions.assertDoesNotThrow(() ->
kafkaLogCollectClient.consume0(Collections.singletonList(shenyuRequestLog)));
+ verify(producer).flush();
+ }
+
+ private void setProducer(final KafkaProducer<String, String> producer)
throws NoSuchFieldException, IllegalAccessException {
+ Field field = KafkaLogCollectClient.class.getDeclaredField("producer");
+ field.setAccessible(true);
+ field.set(kafkaLogCollectClient, producer);
+ }
}