【问题标题】:What is the max number of async threads created for kafkatemplate async response为 kafkatemplate 异步响应创建的最大异步线程数是多少
【发布时间】:2020-08-08 11:05:38
【问题描述】:

“一个 ForkJoinPool 是用给定的目标并行级别构造的;默认情况下,等于可用处理器的数量。”

假设我的 CPU 有 2 个内核。那么,ForkJoinPool 创建的最大线程数是 4?

假设我正在执行一个异步操作,该操作在使用默认 Forkpool 的循环(比如 10k)操作中返回一个未来对象...那么 Forkpool 将创建多少个线程?

List<ListenableFuture<SendResult<String, String>>> cf = new ArrayList<ListenableFuture<SendResult<String, String>>>();

future = kafkaTemplate.send(topicName, message);
cf.add(future);

i++;

future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {

    @Override
    public void onSuccess(SendResult<String, String> result) {
        syso("sent success");
    }

    @Override
    public void onFailure(Throwable ex) {
        System.out.println(" sending failed");
    }
});

而且,在其他课程中,我正在检查所有未来是否都已完成:

    for (ListenableFuture<SendResult<String, String>> m : myFutures) {
        m.get();
    }

【问题讨论】:

  • 你试过了吗?如果线程不可用,该任务将被添加到队列中,你能显示代码吗?
  • 我已经添加了代码.....我正在向 Kakfa 异步发送消息(比如 10k 条消息)并等待未来完成。假设我所有的未来都处于待处理状态......那么将创建多少个线程?
  • 你没有说你从哪里得到这些文档以及它们指的是什么,但你是在问这些文档是否正确?如果您有 2 个 CPU 内核,则将创建 2 个线程。这里的例外是英特尔超线程,操作系统将其视为附加处理器。所以根据你描述的行为,说numThreads = numLogicalCores可能更准确
  • @ChristopherSchneider 我是从那里得到的。但我怀疑我是否正在执行一个异步操作,它会返回一个循环中的未来,如代码中所示,并且它使用默认的 FirkJoinPool 那么可以创建多少个最大线程ForkJoinPool..docs.oracle.com/javase/7/docs/api/java/util/concurrent/…
  • @Deadpool 我在回答中添加了一个示例。

标签: java asynchronous spring-kafka forkjoinpool


【解决方案1】:

没有额外的线程;期货在生产者的 I/O 线程上完成。

这是一个显示回调的测试...

@SpringBootApplication
public class So61415751Application {


    private static final Logger LOG = LoggerFactory.getLogger(So61415751Application.class);


    public static void main(String[] args) {
        SpringApplication.run(So61415751Application.class, args);
    }

    @Bean
    public ApplicationRunner runner(KafkaTemplate<String, String> template) {
        template.setProducerListener(new ProducerListener<String, String>() {
            @Override
            public void onSuccess(ProducerRecord<String, String> producerRecord, RecordMetadata recordMetadata) {
                LOG.info(recordMetadata.toString());
            }
        });
        return args -> {
            IntStream.range(0, 9).forEach(i -> template.send("so61415751", "foo" + i));
            LOG.info("Sent");
            Thread.sleep(10_000);
        };
    }


    @Bean
    public NewTopic topic() {
        return TopicBuilder.name("so61415751").partitions(1).replicas(1).build();
    }

}
spring.kafka.producer.properties.linger.ms=3000

#logging.level.org.springframework.kafka=debug

logging.level.org.apache.kafka=debug

结果

2020-04-24 17:27:46.282  INFO 96084 --- [           main] com.example.demo.So61415751Application   : Sent

...

3 second linger

...

2020-04-24 17:27:49.299  INFO 96084 --- [ad | producer-1] com.example.demo.So61415751Application   : so61415751-0@63
2020-04-24 17:27:49.300  INFO 96084 --- [ad | producer-1] com.example.demo.So61415751Application   : so61415751-0@64
2020-04-24 17:27:49.300  INFO 96084 --- [ad | producer-1] com.example.demo.So61415751Application   : so61415751-0@65
2020-04-24 17:27:49.300  INFO 96084 --- [ad | producer-1] com.example.demo.So61415751Application   : so61415751-0@66
2020-04-24 17:27:49.300  INFO 96084 --- [ad | producer-1] com.example.demo.So61415751Application   : so61415751-0@67
2020-04-24 17:27:49.300  INFO 96084 --- [ad | producer-1] com.example.demo.So61415751Application   : so61415751-0@68
2020-04-24 17:27:49.300  INFO 96084 --- [ad | producer-1] com.example.demo.So61415751Application   : so61415751-0@69
2020-04-24 17:27:49.301  INFO 96084 --- [ad | producer-1] com.example.demo.So61415751Application   : so61415751-0@70
2020-04-24 17:27:49.301  INFO 96084 --- [ad | producer-1] com.example.demo.So61415751Application   : so61415751-0@71

(调用ProducerListener的线程也完成了future)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-01-28
    • 1970-01-01
    • 2010-09-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多