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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 21:36:05