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

如何使用Reactor Kafka批量读取Topic消息并发送到REST端点

Reactor Kafka 批量消费N条消息实现方案

核心是利用Reactor原生的buffer系列操作符实现按固定大小攒批,同时建议搭配超时兜底避免消息量少时长期滞留在内存。

代码调整实现

你原有代码可以按如下修改,假设批次大小N可通过配置注入:

@Slf4j
@Service
public class KafkaConsumerService implements CommandLineRunner {
    @Autowired
    @Qualifier("KafkaConsumerTemplate")
    public ReactiveKafkaConsumerTemplate<String, String> kafkaConsumerTemplate;

    // 可配置的批次大小,默认100条
    @Value("${kafka.consume.batch-size:100}")
    private Integer batchSize;
    // 可选:批次最大等待时间,避免消息少的时候一直攒批不发送,默认5秒,单位毫秒
    @Value("${kafka.consume.batch-max-wait:5000}")
    private Long batchMaxWaitMs;

    // 调用REST接口的方法,此处替换为你的实际调用逻辑
    private Mono<Void> postBatchToRestEndpoint(List<String> batchRecords) {
        log.info("开始推送批次,消息条数:{}", batchRecords.size());
        // 示例:用WebClient调用REST接口的实现逻辑
        // return webClient.post()
        //         .uri("/your/rest/endpoint")
        //         .bodyValue(batchRecords)
        //         .retrieve()
        //         .bodyToMono(Void.class);
        return Mono.empty();
    }

    public Flux<Void> consume() {
        return kafkaConsumerTemplate.receiveAutoAck()
                .doOnNext(consumerRecord -> log.info("received key={}, value={} from topic={}, offset={}",
                        consumerRecord.key(),
                        consumerRecord.value(),
                        consumerRecord.topic(),
                        consumerRecord.offset())
                )
                .map(ConsumerRecord::value)
                .doOnNext(metric -> log.debug("successfully consumed {}={}", Metric[].class.getSimpleName(), metric))
                // 核心攒批逻辑:达到batchSize条或者等待batchMaxWaitMs时间,哪个先触发就输出批次
                .bufferTimeout(batchSize, Duration.ofMillis(batchMaxWaitMs))
                // 过滤空批次,避免消息少的时候超时推送空列表
                .filter(batch -> !batch.isEmpty())
                // 拼接推送REST的逻辑,concatMap保证顺序推送,需要并行可改用flatMap
                .concatMap(this::postBatchToRestEndpoint)
                .doOnError(throwable -> log.error("消费或推送过程出错:{}", throwable.getMessage(), throwable));

    }


    @Override
    public void run(String... args) throws Exception {
        consume().subscribe();
    }
}

注意事项

  • 如果你对消息可靠性要求很高,不建议使用receiveAutoAck自动提交偏移量:自动提交会在消息被消费后直接提交偏移量,如果后续REST调用失败,这批消息会丢失。建议改用receive()手动获取偏移量,批次推送成功后再批量提交偏移量,失败时可以触发重试或者死信队列逻辑。
  • 如果需要严格保证批次大小必须是N条,不需要超时兜底,可以直接用.buffer(N)方法,只有攒够N条才会生成批次。
  • 推送REST如果需要调整并发,可将concatMap替换为flatMap,并指定最大并发数,注意调整后消息的处理顺序会被打乱。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 03:27:03