【问题标题】:How to have sync producer batch messages?如何同步生产者批处理消息?
【发布时间】:2014-06-21 02:20:55
【问题描述】:

我们使用的是 Kafka 0.8 异步生产者,但它正在丢弃消息(并且没有来自另一个线程的 aysnc 响应,或者我们可以继续使用异步)。

我们已将batch.num.messages 设置为 500,并且我们的消费者没有改变。我读到batch.num.messages 只适用于异步生产者而不是同步,所以我需要自己批处理。我们正在使用compression.codec=snappy 和我们自己的序列化程序类。

我的问题有两个:

  • 我可以假设我可以只使用我们自己的序列化程序类,然后自己发送消息吗?

  • 我是否需要担心 Kafka 可能使用的任何特殊的快速选项/参数?

【问题讨论】:

  • 问题解决了吗?您关于 asyc 生产者丢失消息的假设是否正确?我面临着类似的情况,看起来生产者正在非常不可预测地丢失消息。正在使用的版本是0.8.1.1。也尝试0.8.2,但我能够重现它。任何批处理设置会导致不丢弃消息?

标签: apache-kafka clj-kafka


【解决方案1】:

是的,这是因为 batch.num.messages 仅控制 async 生产者的行为。这在相关guide on parameters中明确表示:

使用异步模式时一批发送的消息数。生产者将等待,直到准备好发送此数量的消息或达到 queue.buffer.max.ms。

为了对同步生产者进行批处理,您必须发送消息列表:

public void trySend(List<M> messages) {
    List<KeyedMessage<String, M>> keyedMessages = Lists.newArrayListWithExpectedSize(messages.size());
    for (M m : messages) {
        keyedMessages.add(new KeyedMessage<String, M>(topic, m));
    }
    try {
        producer.send(keyedMessages);
    } catch (Exception ex) {
        log.error(ex)
    }
}

请注意,我在这里使用的是kafka.javaapi.producer.Producer

一旦send 被执行,批处理就会被发送。

我可以假设我可以只使用我们自己的序列化程序类,然后自己发送消息吗? 我需要担心 Kafka 可能使用的任何特殊的快速选项/参数吗?

压缩和序列化器都是不影响批处理的正交功能,但实际上应用于单个消息。

注意会有api变化,async/sync api会统一。

【讨论】:

    猜你喜欢
    • 2017-09-07
    • 1970-01-01
    • 2021-04-30
    • 1970-01-01
    • 1970-01-01
    • 2018-01-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多