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

Spring-Kafka Producer消息发送过慢 求助TPS优化方案

问题描述

基于Spring Boot + Jetty的应用,接收HTTP请求后处理数据,再发送至Kafka,压力测试时Kafka消息发送速度极慢,TPS受限于此。

核心业务代码

private KafkaTemplate<String, Object> kafkaTemplate;

List<Map<String,String>> receive_data = parse(httpRequest);
List<Map<String,String>> processed_data = process(receive_data);

processed_data.forEach(data -> {
  kafkaTemplate.send();
});

线程栈信息

"qtp1206051975-94" #94 prio=5 os_prio=0 tid=0x00007fe054034800 nid=0x34c waiting on condition [0x00007fdfeb4f2000]
   java.lang.Thread.State: TIMED_WAITING (parking)
        at sun.misc.Unsafe.park(Native Method)
        - parking to wait for  <0x0000000788869f80> (a java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject)
        at java.util.concurrent.locks.LockSupport.parkNanos(LockSupport.java:215)
        at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await(AbstractQueuedSynchronizer.java:2163)
        at org.apache.kafka.clients.producer.internals.BufferPool.allocate(BufferPool.java:143)
        at org.apache.kafka.clients.producer.internals.RecordAccumulator.append(RecordAccumulator.java:218)
        at org.apache.kafka.clients.producer.KafkaProducer.doSend(KafkaProducer.java:942)
        at org.apache.kafka.clients.producer.KafkaProducer.send(KafkaProducer.java:865)

当前配置

Kafka Producer配置

props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class.getName());
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, CompressionType.ZSTD.name);
props.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE);
props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 60000);
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 100);
props.put(ProducerConfig.LINGER_MS_CONFIG, 10);
props.put(JsonSerializer.ADD_TYPE_INFO_HEADERS, false);

已尝试调整buffer.memory至64M/128M,无明显效果。

Jetty线程配置

server:
  jetty:
    threads:
      max: 400
      min: 20

问题分析

从线程栈可见,Jetty工作线程卡在BufferPool.allocate的等待状态,核心原因:

  1. 单条消息逐个发送,完全未利用Kafka批量发送机制,导致网络IO频繁,缓冲区无法有效复用。
  2. batch_size=100(单位为字节)远小于单条JSON消息大小,批量机制直接失效。
  3. Jetty最大线程数400过高,大量线程同时向Kafka发送消息,瞬间耗尽生产者缓冲区,引发线程阻塞。

优化方案

1. 批量发送消息

替换循环单条发送为批量发送,减少网络请求次数,充分利用Kafka批量优化:

// 转换为批量ProducerRecord
List<ProducerRecord<String, Object>> records = processed_data.stream()
    .map(data -> new ProducerRecord<>("your_topic_name", data))
    .collect(Collectors.toList());

// 异步批量发送,不阻塞请求线程
CompletableFuture.allOf(
    records.stream()
        .map(record -> kafkaTemplate.send(record))
        .toArray(CompletableFuture[]::new)
).exceptionally(ex -> {
    // 处理发送异常
    log.error("Kafka批量发送失败", ex);
    return null;
});

2. 调整Kafka Producer核心参数

  • 增大batch_size:设置为32768(32KB)或65536(64KB),让生产者能积累足够消息再发送:
    props.put(ProducerConfig.BATCH_SIZE_CONFIG, 65536);
    
  • 调整linger_ms:适当延长linger时间(如50ms),给生产者更多时间积累批量,避免延迟过高:
    props.put(ProducerConfig.LINGER_MS_CONFIG, 50);
    
  • 优化max_in_flight_requests_per_connection:设置为5(默认值),允许同时发送多个请求提升吞吐量(若开启幂等需设为1):
    props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
    
  • 可选调整acks:若业务允许,将acks从all改为1,减少集群确认开销;需强一致性则保留all。

3. 异步解耦请求线程与Kafka发送

使用独立线程池处理Kafka发送,避免Jetty工作线程被阻塞:

// 定义异步线程池
@Bean
public Executor kafkaSendExecutor() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(20);
    executor.setMaxPoolSize(50);
    executor.setQueueCapacity(1000);
    executor.setThreadNamePrefix("kafka-send-");
    executor.initialize();
    return executor;
}

// 业务代码中异步执行发送逻辑
@Autowired
private Executor kafkaSendExecutor;

// ...
kafkaSendExecutor.execute(() -> {
    // 批量或单条发送逻辑
    processed_data.forEach(data -> kafkaTemplate.send("your_topic_name", data));
});

4. 优化Jetty线程配置

过高的Jetty线程数会加剧Kafka缓冲区耗尽问题,建议调整为CPU核心数的2-4倍:

server:
  jetty:
    threads:
      max: 80
      min: 20

5. 监控辅助调优

  • 开启Kafka生产者监控,查看record-send-rate、batch-size-avg、bufferpool-wait-time等指标,验证批量机制是否生效。
  • 检查Kafka集群状态:确保broker节点磁盘IO、网络带宽无瓶颈,副本同步延迟在合理范围。

内容的提问来源于stack exchange,提问作者J.A.R.V.I.S

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 13:57:35