Spring Kafka Template生产者性能求助:15000条消息需8分钟
Spring Kafka Template 生产消息慢?这样调优准提速!
兄弟,15000条消息跑8分钟确实有点拉胯了,咱们来针对性调优Spring Kafka Producer的配置和发送方式,分分钟把速度提上去!
一、先把生产者核心配置拉满
这是最关键的一步,默认配置太保守,得改成适合高吞吐量的参数:
acks配置:如果业务能容忍少量数据丢失(比如日志类消息),直接设为"0",生产者发完就不管了,速度最快;如果要保证至少一个broker收到,设为"1",比默认的"all"快很多("all"需要ISR所有副本确认,延迟高)。batch.size:默认16KB,调大到32KB甚至64KB,让生产者攒更多消息再批量发送,减少网络请求次数。linger.ms:设为5-10ms,让生产者等待一小会儿,凑够一批再发,配合batch.size一起用,最大化批量效率。compression.type:启用压缩,比如"snappy"或者"lz4",压缩比高且解压快,能大幅减少网络传输的数据量。buffer.memory:默认32MB,如果消息量很大,可以调大到64MB,保证生产者有足够缓冲区存放待发送的消息,避免阻塞。max.in.flight.requests.per.connection:设为5,允许同一连接同时发送5个请求,提升并发能力。
修改你的ProducerFactory配置,加上这些参数:
@Bean public ProducerFactory<String, GenericRecord> highSpeedAvroProducerFactory( @Qualifier("highSpeedProducerProperties") KafkaProperties properties) { final Map<String, Object> kafkaPropertiesMap = properties.getKafkaPropertiesMap(); // 高吞吐量配置 kafkaPropertiesMap.put(ProducerConfig.ACKS_CONFIG, "1"); kafkaPropertiesMap.put(ProducerConfig.BATCH_SIZE_CONFIG, 32768); // 32KB kafkaPropertiesMap.put(ProducerConfig.LINGER_MS_CONFIG, 5); kafkaPropertiesMap.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "snappy"); kafkaPropertiesMap.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 67108864); // 64MB kafkaPropertiesMap.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5); kafkaPropertiesMap.put(ProducerConfig.RETRIES_CONFIG, 1); // 减少不必要的重试 return new DefaultKafkaProducerFactory<>(kafkaPropertiesMap); }
二、用批量发送代替单条发送
默认单条调用send()会频繁触发网络请求,改成批量发送效率直接翻倍!
比如把15000条消息打包成列表,一次性发送:
// 构建批量消息 List<ProducerRecord<String, GenericRecord>> batchRecords = new ArrayList<>(15000); for (int i = 0; i < 15000; i++) { GenericRecord avroRecord = buildYourAvroRecord(i); // 你的Avro消息构建逻辑 batchRecords.add(new ProducerRecord<>("your-target-topic", String.valueOf(i), avroRecord)); } // 异步批量发送,不要阻塞等待结果 kafkaTemplate.send(batchRecords) .addCallback( result -> System.out.println("批量发送完成,成功投递" + result.getRecordMetadata().count() + "条"), ex -> System.err.println("批量发送失败:" + ex.getMessage()) );
三、别同步阻塞,用异步发送
如果你之前是每次调用send()后都get()等待结果,那速度慢是必然的!比如这种写法:
kafkaTemplate.send(record).get();
这种同步等待会让生产者变成串行发送,完全浪费了Kafka的并发能力。改成上面的异步回调方式,让生产者同时处理多个发送请求。
四、其他小细节优化
- Avro序列化优化:确保你用的
KafkaAvroSerializer开启了schema缓存,避免每次序列化都去schema registry拉取schema,增加延迟。可以配置schema.registry.url和auto.register.schemas=false(如果schema已经注册过)。 - 检查Kafka集群状态:如果broker的磁盘IO慢、分区数太少(比如只有1个分区),或者网络带宽不够,也会限制生产者速度。可以给topic多建几个分区,让生产者并行发送到不同分区。
- 升级依赖版本:用最新的Spring Kafka和Kafka客户端版本,新版本会修复很多性能问题,比如Kafka 2.8+的生产者有不少吞吐量优化。
按照这些方法调完,15000条消息应该能在几十秒内发送完成,甚至更快!
内容的提问来源于stack exchange,提问作者Prabhakar D
相关产品推荐
相关产品推荐

