【发布时间】:2019-05-01 16:39:31
【问题描述】:
目前,当我创建生产者来发送我的记录时,例如由于某些原因 kafka 不可用,生产者会无限期地发送相同的消息。例如,在收到此错误 3 次后如何停止生成消息:
Connection to node -1 could not be established. Broker may not be available.
我正在使用 reactor kafka 生产者:
@Bean
public KafkaSender<String, String> createSender() {
return KafkaSender.create(senderOptions());
}
private SenderOptions<String, String> senderOptions() {
Map<String, Object> props = new HashMap<>();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getBootstrapServers());
props.put(ProducerConfig.CLIENT_ID_CONFIG, kafkaProperties.getClientId());
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.RETRIES_CONFIG, kafkaProperties.getProducerRetries());
return SenderOptions.create(props);
}
然后用它来发送记录:
sender.send(Mono.just(SenderRecord.create(new ProducerRecord<>(topicName, null, message), message)))
.flatMap(result -> {
if (result.exception() != null) {
return Flux.just(ResponseEntity.badRequest()
.body(result.exception().getMessage()));
}
return Flux.just(ResponseEntity.ok().build());
})
.next();
【问题讨论】:
-
我们可以提供一些您如何生成记录的代码吗?另外,请分享更多堆栈跟踪
-
更新了我的帖子。
-
不是说你的`props.put(ProducerConfig.RETRIES_CONFIG, kafkaProperties.getProducerRetries());`对producer有影响吗?关于此事的堆栈跟踪怎么样?或者至少更多的日志......
-
我在日志中不断看到以下消息:WARN 7468 --- [client] org.apache.kafka.clients.NetworkClient : [Producer clientId=mycliet] 无法建立到节点 -1 的连接。经纪人可能不可用。就这样。该属性不会改变任何东西
-
我们有没有机会在 GitHub 上有一个简单的项目来玩?
标签: spring apache-kafka kafka-producer-api spring-kafka