【问题标题】:Batching not working in Kafka Producer with Scala使用 Scala 在 Kafka Producer 中批处理不起作用
【发布时间】:2018-01-17 23:18:39
【问题描述】:

我正在用 Scala 编写 Producer,我想做批处理。批处理应该工作的方式是,它应该将消息保存在队列中直到它已满,然后将所有消息一起发布到主题上。但不知何故,它不起作用。在我开始发送消息的那一刻,它开始一一发布消息。有谁知道如何在 Kafka Producer 中使用批处理。

val kafkaStringSerializer = "org.apache.kafka.common.serialization.StringSerializer"
      val batchSize: java.lang.Integer = 163840
      val props = new Properties()
      props.put("key.serializer", kafkaStringSerializer)
      props.put("value.serializer", kafkaStringSerializer)
      props.put("batch.size", batchSize);
      props.put("bootstrap.servers", "localhost:9092")

      val producer = new KafkaProducer[String,String](props)

      val TOPIC="topic"
      val inlineMessage = "adsdasdddddssssssssssss"

      for(i<- 1 to 10){
        val record: ProducerRecord[String, String] = new ProducerRecord(TOPIC, inlineMessage )
        val futureResponse: Future[RecordMetadata] =  producer.send(record)
        futureResponse.isDone
        println("Future Response ==========>" + futureResponse.get().serializedValueSize())
      }

【问题讨论】:

    标签: scala apache-kafka kafka-producer-api


    【解决方案1】:

    你必须在你的道具中设置linger.ms

    默认情况下,它为零,这意味着如果可能,消息会立即发送。 您可以增加它(例如 100)以便进行批处理 - 这意味着更高的延迟,但更高的吞吐量。

    batch.size 是一个最大值:如果你在linger.ms 过去之前到达它,数据将被发送而无需等待更多时间。

    要查看实际发送的批次,您需要配置您的日志记录(批处理在后台线程上完成,您将无法查看使用生产者 api 完成的批次 - 您无法发送或接收批次,只发送一条记录并接收它的响应,通过批处理与代理的通信是在内部完成的)

    首先,如果尚未完成,请绑定一个 log4j 属性文件 (Dlog4j.configuration=file:path/to/log4j.properties)

    log4j.rootLogger=WARN, stderr
    log4j.logger.org.apache.kafka.clients.producer.internals.Sender=TRACE, stderr
    
    log4j.appender.stderr=org.apache.log4j.ConsoleAppender
    log4j.appender.stderr.layout=org.apache.log4j.PatternLayout
    log4j.appender.stderr.layout.ConversionPattern=[%d] %p %m (%c)%n
    log4j.appender.stderr.Target=System.err
    

    例如,我会收到

    TRACE Sent produce request to 2: (type=ProduceRequest, magic=1, acks=1, timeout=30000, partitionRecords=({test-1=[(record=LegacyRecordBatch(offset=0, Record(magic=1, attributes=0, compression=NONE, crc=2237306008, CreateTime=1502444105996, key=0 bytes, value=2 bytes))), (record=LegacyRecordBatch(offset=1, Record(magic=1, attributes=0, compression=NONE, crc=3259548815, CreateTime=1502444106029, key=0 bytes, value=2 bytes)))]}), transactionalId='' (org.apache.kafka.clients.producer.internals.Sender)
    

    这是一组 2 个数据。批处理将包含发送到同一代理的记录

    然后,使用 batch.size 和 linger.ms 来看看区别。请注意,一条记录包含一些开销,因此 1000 的 batch.size 不会包含 10 条大小为 100 的消息

    请注意,我没有找到说明所有记录器及其功能的文档(例如 log4j.logger.org.apache.kafka.clients.producer.internals.Sender)。你可以在rootLogger上开启DEBUG/TRACE,找到你想要的数据,或者explore the code

    【讨论】:

    • 我做到了。我有 props.put("linger.ms", 5000)。但仍然无法正常工作。现在我看到我的消息延迟了 5 秒。消息仍在一条一条地传来,但延迟了 5 秒。
    • 消息是分批存储和分批获取的,但它们仍然作为单独的消息呈现给消费者。不要期望阅读一条消息并获得一批消息作为响应。这不是 Kafka 批处理的工作原理。
    • 所以,如果我必须验证我的批处理是否正确完成。如何验证?
    • 您可以查看代理上的 JMX 指标,该指标计算生产请求的数量和获取请求的数量。如果有批处理,那么请求的数量将低于消息的数量。另一种方法是使用wireshark 或其他网络数据包分析器来查看在线上的Kafka 请求和内容。
    【解决方案2】:

    您正在同步向 Kafka 服务器生成数据。也就是说,当你用futureResponse.get 调用producer.send 时,它只会在数据存储到Kafka Server 后返回。

    将响应存储在单独的列表中,并在for 循环之外调用futureResponse.get

    默认configuration,Kafka支持批处理,见linger.msbatch.size

    List<Future<RecordMetadata>> responses = new ArrayList<>();
    for (int i=1; i<=10; i++) {
        ProducerRecord<String, String> record = new ProducerRecord<>(TOPIC, inlineMessage);
        Future<RecordMetadata> response = producer.send(record);
        responses.add(response);
    }
    
    for (Future<RecordMetadata> response : responses) {
        response.get(); // verify whether the message is sent to the broker.
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-12-19
      • 1970-01-01
      • 2011-08-03
      • 2020-04-05
      • 2020-07-31
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多