我有下面提到的不同模型,我创建Kafka生成器并调用不同的方法,但不确定什么是正确的方法来对其进行编程,以便流不会中断,性能不会受到影响。请帮帮忙。
模型1:
for(int i=1; i < 100; i++){
Producer<String, String> producer = new KafkaProducer<String, String>(props);
ProducerRecord<String, String> data = new ProducerRecord<String, String>(
topicName,
String.valueOf(i)
);
producer.send(data);
producer.close();
}模型2:
Producer<String, String> producer = new KafkaProducer<String, String>(props);
for(int i=1; i < 100; i++) {
ProducerRecord<String, String> data = new ProducerRecord<String, String>(
topicName,
String.valueOf(i)
);
producer.send(data);
producer.close();
}模型3:
Producer<String, String> producer = new KafkaProducer<String, String>(props);
for(int i=1; i < 100; i++){
ProducerRecord<String, String> data = new ProducerRecord<String, String>(
topicName,
String.valueOf(i)
);
producer.send(data);
}
producer.close();模型4:
for(int i=1; i < 100; i++){
Producer<String, String> producer = new KafkaProducer<String, String>(props);
ProducerRecord<String, String> data = new ProducerRecord<String, String>(
topicName,
String.valueOf(i)
);
producer.send(data);
producer.flush();
producer.close();
}模型5:
Producer<String, String> producer = new KafkaProducer<String, String>(props);
for(int i=1; i < 100; i++){
ProducerRecord<String, String> data = new ProducerRecord<String, String>(
topicName,
String.valueOf(i)
);
producer.send(data);
producer.flush();
producer.close();
}模型6:
Producer<String, String> producer = new KafkaProducer<String, String>(props);
for(int i=1; i < 100; i++){
ProducerRecord<String, String> data = new ProducerRecord<String, String>(
topicName,
String.valueOf(i)
);
producer.send(data);
producer.flush();
}
producer.close();发布于 2017-12-15 21:15:18
通过以下更改,模型3似乎应该是正确的
Producer<String, String> producer = new KafkaProducer<String, String>(props);
try {
for (int i = 1; i < 100; i++) {
ProducerRecord<String, String> data = new ProducerRecord<String, String>(topicName, String.valueOf(i));
producer.send(data);
}
} catch (Exception e) {
e.printStackTrace();
} finally {
producer.close();
}发布于 2017-12-15 21:27:40
您可以使用以下示例:
Properties props = new Properties();
props.put("batch.size", 16384);
props.put("buffer.memory", 33554432);
Producer<String, String> producer = new KafkaProducer<>(props);
for (int i = 0; i < 100; i++)
producer.send(new ProducerRecord<String, String>("my-topic", Integer.toString(i), Integer.toString(i)));
producer.close();正如您所知,send()方法是异步的。当被调用时,它会将记录添加到挂起记录发送的缓冲区中,并立即返回。这使得生产者可以批量处理单个记录,以提高效率。
我们可以将buffer.memory或batch.size设置为自动刷新
https://stackoverflow.com/questions/47832944
复制相似问题