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
相关产品推荐
相关产品推荐

