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 f9ff52ea0e fix logging-kafka : avoid sending a record during
initialization. (#7150)
f9ff52ea0e is described below
commit f9ff52ea0e7f803c410ef496728bb1938aaebb2c
Author: Jerry聊AI <[email protected]>
AuthorDate: Mon Sep 21 19:12:17 2026 +0800
fix logging-kafka : avoid sending a record during initialization. (#7150)
Co-authored-by: aias00 <[email protected]>
---
.../kafka/client/KafkaLogCollectClient.java | 27 +++++++------------
.../kafka/kafka/KafkaLogCollectClientTest.java | 31 +++++++++++++++-------
2 files changed, 30 insertions(+), 28 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 3cbb1ecbab..e919111dab 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
@@ -26,9 +26,6 @@ import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.KafkaException;
import org.apache.kafka.common.config.SaslConfigs;
-import org.apache.kafka.common.errors.AuthorizationException;
-import org.apache.kafka.common.errors.OutOfOrderSequenceException;
-import org.apache.kafka.common.errors.ProducerFencedException;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.common.utils.JsonUtils;
@@ -96,25 +93,19 @@ public class KafkaLogCollectClient extends
AbstractLogConsumeClient<KafkaLogColl
.format("org.apache.kafka.common.security.scram.ScramLoginModule required
username=\"{0}\" password=\"{1}\";",
config.getUserName(),
config.getPassWord()));
}
- producer = new KafkaProducer<>(props);
- ProducerRecord<String, String> record = new
ProducerRecord<>(this.topic, StringSerializer.class.getName(),
StringSerializer.class.getName());
try {
- producer.send(record);
- LOG.info("init kafkaLogCollectClient success");
- } catch (ProducerFencedException | OutOfOrderSequenceException |
AuthorizationException e) {
- // We can't recover from these exceptions, so our only option is
to close the producer and exit.
- LOG.error("Init kafkaLogCollectClient error, We can't recover from
these exceptions, so our only option is to close the producer and exit", e);
- producer.close();
- return false;
+ producer = new KafkaProducer<>(props);
+ producer.partitionsFor(this.topic);
+ LOG.info("kafka topic metadata fetched successfully");
+ return true;
} catch (KafkaException e) {
- // For all other exceptions, just abort the transaction and try
again.
- LOG.error(
- "init kafkaLogCollectClient error,Exceptions other than
ProducerFencedException or OutOfOrderSequenceException or
AuthorizationException"
- + ", just abort the transaction and try again", e);
- producer.close();
+ LOG.error("Failed to initialize kafka producer", e);
+ if (Objects.nonNull(producer)) {
+ producer.close();
+ producer = null;
+ }
return false;
}
- return true;
}
/**
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 e0b5cd6b83..a307fc3e4b 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.errors.TimeoutException;
import org.apache.shenyu.common.dto.PluginData;
import org.apache.shenyu.common.utils.GsonUtils;
import org.apache.shenyu.plugin.logging.common.entity.ShenyuRequestLog;
@@ -25,13 +26,14 @@ import
org.apache.shenyu.plugin.logging.kafka.client.KafkaLogCollectClient;
import org.apache.shenyu.plugin.logging.kafka.config.KafkaLogCollectConfig;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
-import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Test;
import org.mockito.MockedConstruction;
-import java.lang.reflect.Field;
-
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
/**
* The Test Case For KafkaLogCollectClient.
@@ -58,14 +60,23 @@ public class KafkaLogCollectClientTest {
}
@Test
- @Disabled
- public void testInitClient() throws NoSuchFieldException,
IllegalAccessException {
+ public void testInitClientDoesNotSendRecord() {
try (MockedConstruction<KafkaProducer> construction =
mockConstruction(KafkaProducer.class)) {
- kafkaLogCollectClient.initClient(globalLogConfig);
- Field field =
kafkaLogCollectClient.getClass().getDeclaredField("topic");
- field.setAccessible(true);
- Assertions.assertEquals(field.get(kafkaLogCollectClient),
"shenyu-access-logging");
- kafkaLogCollectClient.close();
+
Assertions.assertTrue(kafkaLogCollectClient.initClient0(globalLogConfig));
+ Assertions.assertEquals(1, construction.constructed().size());
+
verify(construction.constructed().get(0)).partitionsFor("shenyu-access-logging");
+ verify(construction.constructed().get(0), never()).send(any());
+ }
+ kafkaLogCollectClient.close0();
+ }
+
+ @Test
+ public void testInitClientFailsWhenTopicMetadataIsUnavailable() {
+ try (MockedConstruction<KafkaProducer> construction =
mockConstruction(KafkaProducer.class,
+ (mock, context) ->
when(mock.partitionsFor("shenyu-access-logging")).thenThrow(new
TimeoutException("metadata unavailable")))) {
+
Assertions.assertFalse(kafkaLogCollectClient.initClient0(globalLogConfig));
+ verify(construction.constructed().get(0)).close();
+ verify(construction.constructed().get(0), never()).send(any());
}
}
}