Aias00 commented on code in PR #7150:
URL: https://github.com/apache/shenyu/pull/7150#discussion_r4059874711


##########
shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-kafka/src/main/java/org/apache/shenyu/plugin/logging/kafka/client/KafkaLogCollectClient.java:
##########
@@ -96,25 +93,19 @@ public boolean initClient0(@NonNull final 
KafkaLogCollectConfig.KafkaLogConfig c
                             
.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);

Review Comment:
   Suggestion (non-blocking): this call blocks until metadata arrives or 
`max.block.ms` elapses (60 s by default), and `initClient0` is reached from 
`LoggingKafkaPluginDataHandler#doRefreshConfig`, i.e. the config-sync path - an 
unreachable broker would park it for up to a minute per refresh.
   
   The previous `send()` was not better here (`KafkaProducer#send` calls 
`waitOnMetadata` with the same budget), so this is not a regression, just worth 
bounding now that fetching metadata is the whole point of the probe. Since 
kafka-clients 3.9.2 has no `partitionsFor(String, Duration)` overload (verified 
with javap against the exact version this reactor depends on), the only lever 
is:
   
       props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 5_000L);
   
   Also worth considering: when the topic does not exist and auto-creation is 
disabled, `partitionsFor` simply returns an empty list instead of raising, so a 
mistyped topic still reports a successful init.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to