【发布时间】:2023-03-24 10:49:01
【问题描述】:
我有一个有趣的时刻,来自 org.apache.kafka.clients 包的 KafkaProducer 类。 根据文档,flush() 方法的执行应该阻止代码执行,直到所有消息都发送完毕。
我有这样的代码:
messageSender.flush();
messageSender.close();
方法 messageSender.flush() 为我所有的生产者执行刷新:
public void flush() {
producers.forEach(Producer::flush);
}
在执行第一个代码块之前,我通过 send() 方法发送了一些消息。 但在结束之后,我看到,并非所有消息都在生产者关闭之前发送。
如果我将第一个代码块更改为:
messageSender.flush();
try {
Thread.sleep(500);
} catch (InterruptedException e) {
e.printStackTrace();
}
messageSender.close();
所有消息都已发送。
我做错了什么?
我的制作人有这样的配置:
(ProducerConfig.BATCH_SIZE_CONFIG, "65536");
(ProducerConfig.ACKS_CONFIG, "0");
(ProducerConfig.LINGER_MS_CONFIG, "40");
(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
【问题讨论】:
-
您不需要多个生产者实例。它们是线程安全的
-
我知道,它们是线程安全的。但是一个实例的吞吐量对于我的需求来说太低了。
标签: java apache-kafka kafka-producer-api