【问题标题】:Do newer versions of Kafka producers still have "producer.type"?较新版本的 Kafka 生产者是否仍然具有“producer.type”?
【发布时间】:2018-02-03 01:28:54
【问题描述】:

旧版本的文档说这是基本属性之一。

较新版本的文档根本没有提及。

较新版本的 Kafka 生产者是否还有producer.type

或者,新的制作人总是async,我应该打电话给future.get(),让它成为sync

【问题讨论】:

    标签: java apache-kafka producer


    【解决方案1】:

    新的生产者总是异步的,你应该调用 future.get() 让它同步。当像添加future.get()这样简单的东西给你基本相同的功能时,创建两个api方法是不值得的。

    来自 send() here 的文档

    https://kafka.apache.org/0110/javadoc/index.html?org/apache/kafka/clients/producer/KafkaProducer.html

    由于发送调用是异步的,它返回一个 Future 将分配给此记录的 RecordMetadata。调用 get() on 这个未来将阻塞,直到相关的请求完成,然后 返回记录的元数据或抛出任何异常 发送记录时发生。

    如果你想模拟一个简单的阻塞调用,你可以调用 get() 立即方法:

    byte[] key = "key".getBytes();
    byte[] value = "value".getBytes();  
    ProducerRecord<byte[],byte[]> record = new ProducerRecord<byte[],byte[]>("my-topic", key, value);
    producer.send(record).get();
    

    【讨论】:

    【解决方案2】:

    为什么要让send() 同步?

    这是一个 kafka 功能,用于批处理消息以获得更好的吞吐量。

    异步发送

    批处理是效率的主要驱动力之一,为了启用批处理,Kafka 生产者将尝试在内存中累积数据并在单个请求中发送更大的批处理。批处理可以配置为累积不超过固定数量的消息,并且等待不超过某个固定延迟限制(例如 64k 或 10 ms)。这允许累积更多要发送的字节,并在服务器上进行少量较大的 I/O 操作。这种缓冲是可配置的,并提供了一种机制来权衡少量的额外延迟以获得更好的吞吐量。

    由于 api 仅支持异步方法,因此无法进行发送同步,但是您可以指定一些配置来做一些工作。

    您可以将 batch.size 设置为 0。在这种情况下,消息 bacthing 被禁用。

    但我认为您应该保留 batch.size 默认值并将 linger.ms 设置为 0(这也是默认值)。在这种情况下,如果多条消息同时进来,它们将立即批量发送。

    生产者将在请求传输之间到达的所有记录组合成一个批处理请求。通常,这仅在记录到达速度快于发送速度时才会在负载下发生。

    如果您想确保消息成功发送和持久化,您可以将 acks 设置为 -1 或 1 并将 retries 设置为 3(例如)

    更多生产者配置信息,可以参考https://kafka.apache.org/documentation/#producerconfigs

    【讨论】:

    • 我相信同步发送在某些情况下很有用。或者他们过去为什么提供它?
    • 例如,我想确保我的消息被持久化
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-03
    • 2010-12-10
    • 1970-01-01
    • 2010-10-24
    • 2012-09-01
    • 2017-10-18
    相关资源
    最近更新 更多