如何使用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
相关产品推荐
相关产品推荐

