【发布时间】: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