【问题标题】:Collecting messages which are not successfully written on Kafka收集未成功写入 Kafka 的消息
【发布时间】:2017-02-01 21:48:12
【问题描述】:

我正在读取一个文件并将每条记录转储到 Kafka 上。这是我的生产者代码:

public void produce(String topicName, String filePath, String bootstrapServers, String encoding) {
     try (BufferedReader bf = getBufferedReader(filePath, encoding);
                 KafkaProducer<Object, String> producer = initKafkaProducer(bootstrapServers)) {
                String line;
                long count = 0;
                while ((line = bf.readLine()) != null) {
                    count++;
                    producer.send(new ProducerRecord<>(topicName, line), (metadata, e) -> {
                        if(e != null){
                            e.printStackTrace();
                            //write record to some file.
                        }
                    });
                }
                producer.flush();
                CustomLogger.log("Done producing data messages. Total no of records produced:" + count);
            } catch (IOException e) {
                Throwables.propagate(e);
            }
}
 private static KafkaProducer<Object, String> initKafkaProducer(String bootstrapServer) {
        Properties properties = new Properties();
        properties.put("bootstrap.servers", bootstrapServer);
        properties.put("key.serializer", StringSerializer.class.getCanonicalName());
        properties.put("value.serializer", StringSerializer.class.getCanonicalName());
        properties.put("acks", "-1");
        properties.put("retries", 4);
        return new KafkaProducer<>(properties);
    }

private BufferedReader getBufferedReader(String filePath, String encoding) throws UnsupportedEncodingException, FileNotFoundException {
    return new BufferedReader(new InputStreamReader(new FileInputStream(filePath), Optional.ofNullable(encoding).orElse("UTF-8")));
}

根据我们的基本测试,生成消息可能会由于 TimeoutException 而失败。但是,根据 official documentation of Callback TimeoutException 是一个可重试的异常。意味着在下次重试时可能会产生此消息。因此,如果我在回调中找到 TimeoutException,我不能认为记录发送失败。有没有什么可行的方法可以肯定地说记录发送失败并将其记录在单独的文件中?

【问题讨论】:

    标签: java apache-kafka kafka-producer-api


    【解决方案1】:

    我简要地查看了代码,并不认为您需要在这里区分可重试和不可重试异常,因为这已经在 KafkaProducer 中发生了。

    当您使用大于 1 的 retries 值配置生产者时,它会重新发送任何失败并出现可重试异常的消息(批次),次数与您告诉它的一样多,然后再返回你例外。

    所以基本上,除了生产者放弃的例外,您收到的任何消息。

    查看代码中的completeBatch & canRetry 以确认我的理解,但我个人认为这种行为是有道理的。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-04-24
      • 2019-07-26
      • 2021-06-06
      • 2016-05-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-08-28
      相关资源
      最近更新 更多