You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Java HTTP与Kafka性能测试:while循环是否拖慢性能及优化方案

Spring Reactive + Kafka 性能优化方案

一、当前环境与测试结果

  • 硬件配置:Intel i9-10900K(10核20线程、5GHz)、64GB内存、单块970 Evo NVMe磁盘,所有请求均为本机内部调用
  • 测试数据:
    • HTTP/Tomcat调用:每秒约65次
    • Kafka:最高每秒处理50000条请求,但存在0.5秒延迟

二、当前核心代码与配置

1. HttpClient连接池配置

public ConnectionProvider getConnectionProvider(){
    int maxConnections = 20;
    ConnectionProvider connProvider = ConnectionProvider
    .builder("webclient-conn-pool")
    .maxConnections(maxConnections)
    .maxIdleTime(Duration.of(20, ChronoUnit.SECONDS))
    .maxLifeTime(Duration.of(1000, ChronoUnit.SECONDS))
    .pendingAcquireMaxCount(2000)
    .pendingAcquireTimeout(Duration.ofMillis(20000))
    .build();
    return connProvider;
}

public HttpClient getHttpClient(){
    return HttpClient
    .create(getConnectionProvider());
    //.secure(sslContextSpec -> sslContextSpec.sslContext(webClientSslHelper.getSslContext()))
    /*.tcpConfiguration(tcpClient -> {
            LoopResources loop = LoopResources.create("webclient-event-loop",
                        selectorThreadCount, workerThreadCount, Boolean.TRUE);

            return tcpClient
                    .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 10000)
                   .option(ChannelOption.TCP_NODELAY, true);*/
}

2. Kafka相关配置

application.properties

spring.kafka.bootstrap-servers=PLAINTEXT://localhost:9092,PLAINTEXT://localhost:9093
host.name=localhost

Kafka主题与Admin配置

@Bean
public KafkaAdmin kafkaAdmin() {
    String bootstrapAddress = "localhost:9092";

    Map<String, Object> configs = new HashMap<>();
    configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
    return new KafkaAdmin(configs);
}

@Bean
public NewTopic testTopic() {
    return new NewTopic("test-topic", 6, (short) 1);
}

消息发送定时任务

@Async
@Scheduled(fixedDelay = 15000)
public void scheduleTaskUsingCronExpression() {
   generateCalls();
}

private void generateCalls() {
    try{

        int i = 0;
        System.out.println("start");
        long startTime = System.currentTimeMillis();
        while(i <= 5000){
            //Thread.sleep(5);
            String message = "Test Message sadg sad-";
            kafkaTemplate.send(TOPIC, message + i);
            i++;
        }
        long endTime = System.currentTimeMillis();
        System.out.println((endTime - startTime));
        System.out.println("done");
    }
    catch(Exception e){
        e.printStackTrace();
    }
    System.out.println("RUNNING");
}

Kafka消费者代码

@KafkaListener(topics = "test-topic", groupId = "one", concurrency = "6" )
public void listenGroupFoo(String message) {
    if(message.indexOf("-0") != -1){
        startTime = new Date().getTime();
        System.out.println("Starting Message in group foo: " + message);
    }
    else if(message.indexOf("-100000") != -1){
        endTime = new Date().getTime();
        System.out.println("Received Message in group foo: " + message);
        System.out.println(endTime - startTime);
    }
}

三、性能优化方案分析

针对你提出的7个猜想,逐一给出落地优化建议:

1. 多线程优化

  • HTTP请求:当前while循环同步发送请求,没有利用Spring Reactive的异步特性。改成用Flux.range生成请求,通过flatMap并发执行,同时配合调整后的连接池,提升并发量。
  • Kafka发送:当前循环同步调用kafkaTemplate.send,可改用批量发送——积累N条消息后调用kafkaTemplate.sendBatch,或者用Flux异步批量提交,减少网络交互次数。

2. 修改Tomcat等服务器配置

  • Tomcat端:HTTP吞吐量低的核心瓶颈之一是Tomcat连接数限制,调整以下参数:
    • 增大maxConnections到2000+(默认1000)
    • 设置maxThreads为40-80(对应CPU线程数)
    • 确认启用NIO模式(Spring Boot默认已开启)
    • 调大acceptCount到1000,提升等待队列容量
  • Reactive客户端:当前HttpClient连接池maxConnections仅设20,直接调到200-500,同时启用你注释掉的TCP_NODELAY配置,减少TCP延迟。

3. KafkaTemplate实例优化

KafkaTemplate本身是线程安全的,不需要创建多个实例,但要优化生产者配置:

  • 配置batch.size=16384、linger.ms=5,让客户端批量攒消息后发送,减少网络请求
  • 增大buffer.memory到64MB+,避免消息积压
  • 保持Autowired注入方式即可,无需替换

4. 增加多块磁盘?

当前单NVMe磁盘的顺序写入能力远超5万条短消息的需求,磁盘IO不是瓶颈,暂时不需要增加磁盘,优先优化软件配置。

5. 优化HTTP接收端连接能力

  • 若用Tomcat作为接收端,按上述Tomcat参数调整即可;若用Spring Reactive自带的Netty服务器,调整server.netty.worker-count为CPU核心数*2,server.netty.connection-timeout设为合理值(比如10s)
  • 考虑将Tomcat替换为Netty服务器,配合Reactive模式进一步提升并发处理能力

6. 调整定时任务发送逻辑

  • 当前定时任务瞬间发送5001条消息,导致资源瞬时打满,延迟升高。改成匀速发送:用Flux.interval每秒发送指定数量的消息,或者异步批量提交,让消息生产更均匀
  • 自定义@Async的线程池(核心线程数10、最大20),避免默认SimpleAsyncTaskExecutor创建大量线程带来的上下文切换开销

7. 其他优化手段

  • HTTP请求:
    • 启用HTTP/2,减少TCP连接开销
    • 关闭不必要的日志,降低IO开销
    • 用exchange替代retrieve,设置responseTimeout为0(即发即弃无需等待响应)
  • Kafka:
    • 启用生产者compression.type=snappy,压缩消息大小提升吞吐量
    • 消费者调整fetch.min.bytes和fetch.max.wait.ms,批量拉取消息提升处理效率
  • JVM优化:设置堆内存-Xms32G -Xmx32G,启用G1垃圾收集器,减少GC停顿
  • 系统参数:Linux环境下调整net.core.somaxconn、net.ipv4.tcp_max_syn_backlog增大TCP连接队列,设置vm.swappiness=1减少内存交换

内容的提问来源于stack exchange,提问作者Kevin Nielsen

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 16:05:34