【发布时间】: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