Spring Batch应用升级Spring Boot至2.7.6后性能大幅下降
问题分析与解决方案
Spring Batch 4.3.2(对应Spring Boot 2.4.4及以上版本)对KafkaItemWriter做了核心变更:默认开启waitForSendResults(值为true),会同步等待每条消息的Kafka发送结果,这直接导致了你升级后遇到的性能下降问题。以下是几种恢复异步发送模式的可行方案:
方案1:原生配置回退(推荐)
Spring Batch在引入同步逻辑时,保留了异步模式的开关——通过KafkaItemWriterBuilder设置waitForSendResults(false)即可恢复到4.3.1及之前的异步发送行为,无需自定义Writer:
@Bean public KafkaItemWriter<String, MyKafkaObject> getKafkaItemWriter( KafkaTemplate<String, MyKafkaObject> kafkaProducerTemplate) { LOGGER.info("Creating async KafkaItemWriter"); return new KafkaItemWriterBuilder<String, MyKafkaObject>() .kafkaTemplate(kafkaProducerTemplate) .itemKeyMapper(this.kafkaItemKeyMapper) .waitForSendResults(false) // 关键配置:关闭同步等待 .build(); }
该配置下,KafkaItemWriter调用kafkaTemplate.send()后立即返回,不阻塞等待结果;同时在flush()阶段会自动调用Kafka生产者的刷新逻辑,确保缓存的消息被推送到Broker,平衡了性能与数据可靠性。
方案2:优化自定义NoWaitKafkaItemWriter解决丢数风险
如果你坚持使用自定义Writer,必须修正当前的flush()空实现——空逻辑会导致Kafka生产者缓存的消息可能未发送就结束Chunk,一旦Broker故障或应用重启会丢失数据。优化后的实现如下:
public class NoWaitKafkaItemWriter<K, T> extends KafkaItemWriter<K, T> { private static final Logger LOGGER = LogManager.getLogger(NoWaitKafkaItemWriter.class); private final List<ListenableFuture<SendResult<K, T>>> sendFutures = new ArrayList<>(); @Override protected void writeKeyValue(K key, T value) { LOGGER.info("About to write key {} value {}", key, value); ListenableFuture<SendResult<K, T>> future; if (this.delete) { future = this.kafkaTemplate.sendDefault(key, null); } else { future = this.kafkaTemplate.sendDefault(key, value); } // 记录发送Future,用于flush阶段验证结果 sendFutures.add(future); // 可选:添加失败回调,及时记录发送异常 future.addCallback( result -> LOGGER.debug("Successfully sent key {}", key), ex -> LOGGER.error("Failed to send key {}", key, ex) ); } @Override protected void flush() { LOGGER.info("Flushing pending Kafka sends"); // 触发Kafka生产者刷新,将缓存消息推送到Broker this.kafkaTemplate.flush(); // 等待所有发送请求完成(可根据业务需求添加超时时间) for (ListenableFuture<SendResult<K, T>> future : sendFutures) { try { future.get(5, TimeUnit.SECONDS); } catch (InterruptedException | ExecutionException | TimeoutException e) { LOGGER.error("Error waiting for Kafka send result", e); // 根据业务规则处理:抛出异常终止Batch,或记录后继续 throw new BatchRuntimeException("Kafka send failed during flush", e); } } sendFutures.clear(); } }
这个实现既保留了异步发送的性能优势,又通过flush()阶段的验证确保消息被正确推送,降低了数据丢失风险。
方案3:完全自定义ItemWriter
如果需要更灵活的控制(比如批量发送、自定义失败重试逻辑),可以直接实现ItemWriter接口,绕过Spring Batch原生的KafkaItemWriter:
public class AsyncKafkaItemWriter<T> implements ItemWriter<T> { private static final Logger LOGGER = LogManager.getLogger(AsyncKafkaItemWriter.class); private final KafkaTemplate<String, MyKafkaObject> kafkaTemplate; private final Function<T, String> keyMapper; public AsyncKafkaItemWriter(KafkaTemplate<String, MyKafkaObject> kafkaTemplate, Function<T, String> keyMapper) { this.kafkaTemplate = kafkaTemplate; this.keyMapper = keyMapper; } @Override public void write(List<? extends T> items) throws Exception { List<ListenableFuture<SendResult<String, MyKafkaObject>>> futures = new ArrayList<>(); for (T item : items) { String key = keyMapper.apply(item); LOGGER.info("About to write key {} value {}", key, item); futures.add(kafkaTemplate.sendDefault(key, (MyKafkaObject) item)); } // 触发Kafka生产者刷新 kafkaTemplate.flush(); // 等待所有发送完成(可选,根据可靠性要求调整) for (ListenableFuture<SendResult<String, MyKafkaObject>> future : futures) { try { future.get(5, TimeUnit.SECONDS); } catch (Exception e) { LOGGER.error("Kafka send failed", e); throw new BatchRuntimeException("Failed to send items to Kafka", e); } } } }
配置Bean时直接使用自定义实现:
@Bean public ItemWriter<MyKafkaObject> asyncKafkaItemWriter( KafkaTemplate<String, MyKafkaObject> kafkaProducerTemplate) { return new AsyncKafkaItemWriter<>(kafkaProducerTemplate, this.kafkaItemKeyMapper); }
注意事项
- 异步模式下,需配合Kafka生产者配置(如
acks、retries、linger.ms)来平衡性能与数据可靠性; - 如果业务要求严格的Exactly-Once语义,需结合Spring Batch事务与Kafka事务配置,避免消息重复或丢失;
- 无论采用哪种方案,都建议在Chunk结束时触发Kafka生产者的刷新操作,确保缓存消息被及时推送。
内容的提问来源于stack exchange,提问作者mtam7892
相关产品推荐
相关产品推荐

