Project Reactor Kafka消费消息后API调用无法并行问题咨询
问题根因排查
你的代码存在以下4个核心问题,直接导致API调用串行执行:
- 多余的阻塞包装逻辑:你给reactor-kafka的消费流额外套了一层
Flux.create,还在内部调用了blockLast(),这会直接阻塞消费线程,让原本异步非阻塞的消费流程变成同步串行执行,上游只能等前一条消息处理完才会下发下一条。同时代码中还出现变量引用错误:定义了receiver变量接收消费流,实际却订阅了未声明的kafkaFlux变量。 - 并行逻辑写错位置:你把
parallel(10).runOn(Schedulers.parallel())写在了flatMap内部的单个请求Mono上,parallel算子是用来拆分Flux流为多轨并行,对单个Mono调用没有任何效果,完全起不到并行作用。 - 调度器选型错误:
Schedulers.parallel()是为CPU密集型任务设计的,线程数和CPU核心数强绑定,你给Pod配置的CPU只有300m(相当于0.3核),对应parallel调度器最多只会创建1~2个工作线程,根本撑不起10并发的IO调用。 - 缺少连接池配置:你直接用默认配置创建
WebClient,Reactor Netty默认连接池的最大连接数和CPU核心数绑定,低CPU配额下连接池上限极低,就算线程够也会因为没有可用连接导致请求串行排队。
修复方案
1. 清理错误的消费流包装逻辑
直接使用reactor-kafka原生返回的Flux消费流,不需要额外套Flux.create,删除阻塞的blockLast()调用。
2. 修正并行逻辑,选择正确的调度器
使用带并发数参数的flatMap控制并行度,IO密集型请求统一使用Schedulers.boundedElastic()调度器。
3. 配置WebClient连接池
手动调整WebClient底层连接池的最大连接数,至少大于你设置的并行度。
修正后代码示例
// 1. 配置带自定义连接池的WebClient HttpClient httpClient = HttpClient.create() .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 2000) .responseTimeout(Duration.ofSeconds(2)) .doOnConnected(conn -> conn.addHandlerLast(new ReadTimeoutHandler(2, TimeUnit.SECONDS)) .addHandlerLast(new WriteTimeoutHandler(2, TimeUnit.SECONDS)) ) // 连接池最大连接数调整为20,大于你的并行度即可 .poolResources(PoolResources.create("webclient-pool", 20)); WebClient wc = WebClient.builder() .baseUrl("https://abc.com:8443") .clientConnector(new ReactorClientHttpConnector(httpClient)) .build(); // 2. 清理后的Kafka消费逻辑 ReceiverOptions<Integer, String> options = receiverOptions .subscription(Collections.singleton(topic)) .addAssignListener(partitions -> log.debug("onPartitionsAssigned {}", partitions)) .addRevokeListener(partitions -> log.debug("onPartitionsRevoked {}", partitions)); Flux<ReceiverRecord<Integer, String>> messages = kafkaReceiver.create(options).receive(); // 3. 并行调用API逻辑,单Pod并发数设为10,刚好满足10TPS(接口响应1s) messages.flatMap(record -> { // 建议等API调用成功再提交offset,避免调用失败时消息丢失 return wc.get().uri("/actuator") .retrieve() .bodyToMono(String.class) // 调用成功再ack .doOnSuccess(resp -> record.receiverOffset().acknowledge()) // 异常处理避免流中断 .onErrorResume(e -> { log.error("调用接口失败, 消息key:{}", record.key(), e); // 可以根据需要决定失败时是否提交offset,或者转入死信队列 return Mono.empty(); }); }, 10) // flatMap第二个参数指定并发度为10 .subscribeOn(Schedulers.boundedElastic()) .subscribe();
额外优化建议
- 可以根据实际压测结果调整flatMap的并发数,只要不超过下游接口的限流阈值即可
- 如果需要更高的整体TPS,可以适当增加K8s Pod副本数,4个副本即可达到你的目标120TPS要求
- Kafka消费的并发数和分区数绑定,如果要提升消费速度可以增加Kafka Topic的分区数
内容的提问来源于stack exchange,提问作者OpenSourceDev
相关产品推荐
相关产品推荐

