首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >在创建Kafka生产者并调用send()、flush()和close()方法时,正确的顺序是什么?

在创建Kafka生产者并调用send()、flush()和close()方法时,正确的顺序是什么?
EN

Stack Overflow用户
提问于 2017-12-15 21:01:17
回答 2查看 520关注 0票数 0

我有下面提到的不同模型,我创建Kafka生成器并调用不同的方法,但不确定什么是正确的方法来对其进行编程,以便流不会中断,性能不会受到影响。请帮帮忙。

模型1:

代码语言:javascript
复制
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:

代码语言:javascript
复制
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:

代码语言:javascript
复制
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:

代码语言:javascript
复制
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:

代码语言:javascript
复制
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:

代码语言:javascript
复制
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();
EN

回答 2

Stack Overflow用户

发布于 2017-12-15 21:15:18

通过以下更改,模型3似乎应该是正确的

代码语言:javascript
复制
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();
    }
票数 2
EN

Stack Overflow用户

发布于 2017-12-15 21:27:40

您可以使用以下示例:

代码语言:javascript
复制
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设置为自动刷新

票数 1
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/47832944

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档