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

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);
}

注意事项

  1. 异步模式下,需配合Kafka生产者配置(如acks、retries、linger.ms)来平衡性能与数据可靠性;
  2. 如果业务要求严格的Exactly-Once语义,需结合Spring Batch事务与Kafka事务配置,避免消息重复或丢失;
  3. 无论采用哪种方案,都建议在Chunk结束时触发Kafka生产者的刷新操作,确保缓存消息被及时推送。

内容的提问来源于stack exchange,提问作者mtam7892

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 12:25:26