【问题标题】:what are the drawnbacks and risks replacing default SimpleAsyncTaskExecutor by own Executor用自己的 Executor 替换默认的 SimpleAsyncTaskExecutor 有什么缺点和风险
【发布时间】:2020-02-27 16:24:23
【问题描述】:

个人知识:我从javacodegeeks 读到:“... SimpleAsyncTaskExecutor 对于玩具项目是可以的,但对于任何比这更大的项目,它有点冒险,因为它不限制并发线程并且不重用线程。所以为了安全起见,我们还将添加一个任务执行器bean...”以及来自baeldung的一个非常简单的示例如何添加我们自己的任务执行器。但是我可以找到任何指导来解释后果和一些值得应用的案例。

个人愿望:我正在努力为我们的微服务日志提供一个企业架构,以便在 Kafka 主题上发布。主要针对我的基于日志的情况,“不限制并发线程而不重用它造成的风险”这句话似乎是合理的。

我在本地桌面上成功运行了波纹管代码,但我想知道我是否正确地提供了自定义任务执行器。

我的问题:考虑到我已经在使用 kafkatempla(即默认情况下同步、单例和线程安全,至少据了解,至少用于生成/发送消息),此配置是否真正朝着正确的方向重用线程和避免在使用 SimpleAsyncTaskExecutor 时意外创建线程?

生产者配置

@EnableAsync
@Configuration
public class KafkaProducerConfig {

    private static final Logger LOGGER = LoggerFactory.getLogger(KafkaProducerConfig.class);

    @Value("${kafka.brokers}")
    private String servers;

    @Bean
    public Executor taskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(2);
        executor.setMaxPoolSize(2);
        executor.setQueueCapacity(500);
        executor.setThreadNamePrefix("KafkaMsgExecutor-");
        executor.initialize();
        return executor;
    }

    @Bean
    public Map<String, Object> producerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, servers);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
        return props;
    }

}

制片人

@Service
public class Producer {

    private static final Logger LOGGER = LoggerFactory.getLogger(Producer.class);

    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    @Async
    public void send(String topic, String message) {
        ListenableFuture<SendResult<String, String>> future = kafkaTemplate.send(topic, message);
        future.addCallback(new ListenableFutureCallback<SendResult<String, String>>() {

            @Override
            public void onSuccess(final SendResult<String, String> message) {
                LOGGER.info("sent message= " + message + " with offset= " + message.getRecordMetadata().offset());
            }

            @Override
            public void onFailure(final Throwable throwable) {
                LOGGER.error("unable to send message= " + message, throwable);
            }
        });
    }
}

用于演示目的:

@SpringBootApplication
public class KafkaDemoApplication  implements CommandLineRunner {

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

    }

    @Autowired
    private Producer p;

    @Override
    public void run(String... strings) throws Exception {
        p.send("test", " qualquer messagem demonstrativa");
    }

}

【问题讨论】:

  • 我猜你必须使用 Spring 而不是 Akka Streams 或 Lagom 等其他响应式框架?
  • 是的,Spring 是我们的微服务基础框架。
  • 这个问题太笼统了;这里不适合。也许它适合codereview.stackexchange.com

标签: java spring multithreading spring-boot apache-kafka


【解决方案1】:

这是SimpleAsyncTaskExecutor的默认实现

protected void doExecute(Runnable task) {
    Thread thread = (this.threadFactory != null ? this.threadFactory.newThread(task) : createThread(task));
    thread.start();
}

为每个任务创建新线程,在 Java 中创建线程并不便宜:(Reference)

线程对象使用大量内存,在大型应用程序中,分配和释放许多线程对象会产生大量内存管理开销。

=> 用这个task executor重复执行任务会对应用性能产生负面影响(而且这个executor默认不限制并发任务数)

这就是为什么建议您使用线程池实现的原因,线程创建开销仍然存在,但由于线程被重用而不是 create-fire-forget 而显着减少。

配置ThreadPoolTaskExecutor时,应根据您的应用负载正确定义两个值得注意的参数:

  1. private int maxPoolSize = Integer.MAX_VALUE;

    这是池中的最大线程数。

  2. private int queueCapacity = Integer.MAX_VALUE;

    这是排队的最大任务数。当队列已满时,默认值可能会导致 OutOfMemory 异常。

使用默认值 (Integer.MAX_VALUE) 可能会导致服务器资源不足/崩溃。

您可以通过增加最大池大小setMaxPoolSize() 的数量来提高吞吐量,为了减少加载增加时的预热,将核心池大小设置为更高的值setCorePoolSize()maxPoolSize - corePoolSize 之间的任意数量的线程将在负载增加)

【讨论】:

  • 我在上面的“@Bean public Executor taskExecutor()...”方向是否正确?关于 maxPoolSize 和 queueCapacity 非常值得您的 cmets 使用。但是,您是否错过了一些额外的配置?您是否在我的上述声明中看到任何奇怪的概念或冲突:“我已经在使用 kafkatemplate(即默认情况下同步、单例和线程安全,至少用于生成/发送消息”?在我的公司,我们今天有 5000 万客户使用我们的手机应用程序消耗多个微服务,我们预计每秒有 150 个请求要求记录每个调用
  • 执行器池与 kafkatemplate 无关。您已经使用了异步版本的 kafkatemplate(非阻塞),因此性能应该很快,150 RPS 很小。看看这个result
  • 感谢您的反馈,即每秒 150 个请求很小。 P/lease 只是假设另一个你认为高的数字。关于“无论如何,执行程序池与kafkatemplate无关”我没有明白你的意思。 kafkatemplate 肯定依赖于执行者,不是吗?我知道 kafkatemplate 异步版本很快,但 Manh 的评论对我来说确实有意义“分配和解除分配许多线程对象会产生显着的内存管理开销......这个执行程序默认情况下不限制并发任务的数量”。看来你看到我走错了方向。能说清楚点吗?
  • 根据stackoverflow.com/a/48145863/4148175“生产者的实现是异步的。消息存储在内部队列中等待内部线程发送,这将通过潜在的批处理提高效率”。这样的内线不就是executor提供的吗?如果是这样,拥有一个池我相信它会在一些负面情况下有所帮助(例如,Microsevice 容器和 Kafka 容器之间的连接较低,消息中出现一些意外的繁荣等等。当我说“这样​​的内部线程正是一个由执行器池提供”?
猜你喜欢
  • 2011-12-30
  • 2011-11-07
  • 1970-01-01
  • 2011-11-09
  • 1970-01-01
  • 2011-07-02
  • 1970-01-01
  • 2018-07-17
相关资源
最近更新 更多