【发布时间】:2018-02-03 01:28:54
【问题描述】:
旧版本的文档说这是基本属性之一。
较新版本的文档根本没有提及。
较新版本的 Kafka 生产者是否还有producer.type?
或者,新的制作人总是async,我应该打电话给future.get(),让它成为sync?
【问题讨论】:
标签: java apache-kafka producer
旧版本的文档说这是基本属性之一。
较新版本的文档根本没有提及。
较新版本的 Kafka 生产者是否还有producer.type?
或者,新的制作人总是async,我应该打电话给future.get(),让它成为sync?
【问题讨论】:
标签: java apache-kafka producer
新的生产者总是异步的,你应该调用 future.get() 让它同步。当像添加future.get()这样简单的东西给你基本相同的功能时,创建两个api方法是不值得的。
来自 send() here 的文档
由于发送调用是异步的,它返回一个 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();
【讨论】:
为什么要让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
【讨论】: