【问题标题】:Kafka keeps producing requests even if broker is down即使代理关闭,卡夫卡也会继续产生请求
【发布时间】: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


【解决方案1】:

恐怕clusterAndWaitTime = waitOnMetadata(record.topic(), record.partition(), maxBlockTimeMs); 不参与重试,默认情况下它会迭代直到maxBlockTimeMs = 60000。您可以通过ProducerConfig.MAX_BLOCK_MS_CONFIG 属性为生产者减少此选项:

public static final String MAX_BLOCK_MS_CONFIG = "max.block.ms";
    private static final String MAX_BLOCK_MS_DOC = "The configuration controls how long <code>KafkaProducer.send()</code> and <code>KafkaProducer.partitionsFor()</code> will block."
                                                    + "These methods can be blocked either because the buffer is full or metadata unavailable."
                                                    + "Blocking in the user-supplied serializers or partitioner will not be counted against this timeout.";

更新

我们可以这样解决问题:

@PostMapping(path = "/v1/{topicName}")
public Mono<ResponseEntity<?>> postData(
    @PathVariable("topicName") String topicName, String message) {
    return sender.send(Mono.just(SenderRecord.create(new ProducerRecord<>(topicName, null, message), message)))
        .flatMap(result -> {
            if (result.exception() != null) {
                sender.close();
                return Flux.just(ResponseEntity.badRequest()
                    .body(result.exception().getMessage()));
            }
            return Flux.just(ResponseEntity.ok().build());
        })
        .next();
}

注意sender.close();以防出错。

我认为是时候针对 Reactor Kafka 项目提出问题,以允许关闭生产者出错。

【讨论】:

  • 所以我尝试了这个选项。对于我的情况,它会在超时后抛出异常:org.apache.kafka.common.errors.TimeoutException: 1000 ms 后更新元数据失败。但是在那之后生产者仍然继续向kafka发送请求....
  • 考虑使用SenderOptions.stopOnError(true)。有关更多信息,请参阅其 Javadocs。
  • 不知道为什么,但它仍然没有帮助......你有没有尝试重现我的问题?
  • 是的,您的测试用例非常不言自明,但它在WebTestClient 中的Mono.blockingGet() 5 秒后存在,所以我看不到任何其他内容。
  • 设置此props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 3000); 后,测试失败并出现预期错误:Caused by: java.lang.AssertionError: Status expected:&lt;200 OK&gt; but was:&lt;400 BAD_REQUEST&gt;。因此,不会对同一记录进行任何重试。
【解决方案2】:

您可以使用circuit breaker pattern 解决此类问题,但在应用此模式之前,请尝试查找根本原因,并且您的 ProducerConfig.RETRIES_CONFIG 属性似乎在某处被覆盖。

【讨论】:

    【解决方案3】:

    而不是专注于错误。解决问题 - 它没有连接到代理

    您没有在撰写文件中覆盖它,因此您的应用正在尝试连接到自身

    bootstrap-servers: ${KAFKA_BOOTSTRAP_URL:localhost:9092} 
    

    在 compose yml 中,您似乎忘记了这个

    rest-proxy:
       environment:
           KAFKA_BOOTSTRAP_URL: kafka:9092
    

    或者,如果可能,您可以使用现有的 Confluent REST 代理 docker 映像,而不是重新发明轮子

    【讨论】:

    • 我想了解我的应用在出现错误时的行为。
    • 是的,您连接到错误的地址...所以,不知道为什么这值得一票否决
    • 据我所知,Confluent 休息代理正在阻塞。我想创建一个不阻塞的。
    • 因为我想知道如何解决这个特定问题
    • Http 请求通常被阻塞
    猜你喜欢
    • 2020-08-27
    • 2017-02-21
    • 1970-01-01
    • 2019-04-09
    • 2021-01-20
    • 2019-09-24
    • 1970-01-01
    • 2016-05-20
    • 2020-02-20
    相关资源
    最近更新 更多