【问题标题】:KafkaProducer flush does not block the main threadKafkaProducer flush 不阻塞主线程
【发布时间】: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


【解决方案1】:

冲洗

你的主线程不会被producer.flush()阻塞太久, 这是因为使用 0 作为ACK,生产者将更像UDP 发送者:如果调用了发送,则消息成功完成,没有任何保证。所以这个配置没有真正的block-until-received机制,flush 清空的缓冲区已经几乎是空的了。


关闭

在发送消息之前应该阻塞主线程的另一个方法是close() 方法。来自文档:

close()

此方法会阻塞,直到所有先前发送的请求完成。


根据这个逻辑,你说得对,代码应该等到所有消息都发送完毕。您的配置中的问题似乎是 ACKS=0 属性。

ACKs

值为 0 时,生产者甚至不会等待来自 经纪人。 它立即认为写入成功 记录发送出去

可能发生的情况是这样的:close() 将等待直到所有请求都完成,但 ACK 为 0,它将在发送成功消息时撒谎。 producer 不会等待确认,因此会假设每条消息都已正确发送; close() 操作将在之后立即执行,对已完成的请求没有任何真正的保证。在正常情况下,这将主要影响到最后发送的请求。


我会尝试的是:

  • ACKs = 1至少

这将使flush()close() 等待更久,因为生产者需要来自代理的确认才能将请求标记为成功。

  • producer.close(500,TimeUnit.MILLISECONDS)

This close method 等待生产者完成发送所有未完成请求的超时。这也将有助于处理这些最后的消息,并且是您手动处理睡眠的“官方”方式。将其设置为您希望的等待时间。

  • 重复使用相同的producer

不仅在src comments中推荐,而且对一般性能也有帮助。

生产者是线程安全的,跨线程共享单个生产者实例通常比拥有多个实例更快。

【讨论】:

  • 非常感谢!由于吞吐量的想法,我有 ack=0 。我对使用来自 Kafka 的数据的服务进行了一种性能测试。那么,您能否纠正我,ack=1 只会增加 close() 方法执行的延迟,或者所有发送通常会变慢?
  • 是的,每个请求的完成都会变慢,而且如果设置了retries,您会注意到为了将消息标记为正确发送,需要对代理进行更多调用。一般来说,这意味着通常的发送流需要一些额外的延迟,以及flush()close() 的完成时间。这更像是从UDP 转到更常见的TCP 发送
猜你喜欢
  • 2018-02-16
  • 1970-01-01
  • 2015-12-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多