【发布时间】:2020-03-12 17:14:41
【问题描述】:
Java 应用程序向 Kafka 集群发送或发布消息。下面给出了发送记录的代码。由于此应用程序向 kafka 发送大量记录消息,因此调试起来有点困难。但是我看到有些消息没有发送到 kafka,因为我在同一主题的 kafka 代理订阅者中看不到其中一些消息。所以看起来这个管道中有一些数据泄漏。
为此,我添加了java.util.concurrent.Future.isDone() 方法来查看响应。它以错误的方式回应。所以我很困惑到底哪里存在泄漏。是否无法将记录发送到kafka?还是在 Kafka 代理中,在将记录放入主题之前无法处理记录?
// setting properties from config xml..
propertiess.put("security.protocol", properties.getProperty("security_protocol"));
properties.put("acks", properties.getProperty("kafka_acks"));
...
...
//Producer initialization
Producer<String, String> producer = new KafkaProducer<>(kprops);
...
...
ProducerRecord<String, String> producerRecord = new ProducerRecord<>(topicName, key, newValue);
Future<RecordMetadata> kafkaResponse = producer.send(producerRecord);
String kafkaSuccessStatus = kafkaResponse.isDone() ? "Sending message to kafka completed" : "Sending message to kafka not completed";
LOGGER.debug(kafkaSuccessStatus);
我正在使用 isDone() 来检查事件是否成功。这是正确的做法吗?如果没有,我们是否可以通过其他方式获得更多信息。由于我在日志中看不到任何错误或信息,因此很难找出到底发生了什么以及数据被丢弃的位置。我用的是kafka-0.11版本。
关于应用程序的信息:这个应用程序每天需要处理非常大量(数十亿)的记录,并且必须发布到 Kafka 代理。这是一个多线程应用程序从文件中读取行并将其发送到 kafka。
【问题讨论】:
-
这会检查他们是否准备好在下一批中发送,我相信。您可以在运行时关闭挂钩上添加
producer.flush()和/或producer.close()
标签: java apache-kafka kafka-producer-api