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);
+    }
 }

Reply via email to