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的等待状态,核心原因:
- 单条消息逐个发送,完全未利用Kafka批量发送机制,导致网络IO频繁,缓冲区无法有效复用。
batch_size=100(单位为字节)远小于单条JSON消息大小,批量机制直接失效。- 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
相关产品推荐
相关产品推荐

