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]