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

GCP Pub/Sub订阅消费极慢问题排查及并发配置疑问

问题

在GCP Pub/Sub的主题/订阅中积压了约10000条消息,当前使用如下代码及并发配置:

并发配置属性

spring.cloud.gcp.pubsub.subscriber.executor-threads=20
spring.cloud.gcp.pubsub.subscriber.parallel-pull-count=4

消费代码

@Bean
public PubSubInboundChannelAdapter messageChannelAdapterTry(
        @Qualifier("tryInputChannel") MessageChannel inputChannel,
        PubSubTemplate pubSubTemplate) {
    PubSubInboundChannelAdapter adapter = new PubSubInboundChannelAdapter(
            pubSubTemplate, tryConfig.getSubscription());
    adapter.setOutputChannel(inputChannel);
    adapter.setPayloadType(tryClass.class);
    return adapter;
}

@ServiceActivator(inputChannel = "tryInputChannel")
public void MessageReceiver(tryClass tryMessage, @Header(GcpPubSubHeaders.ORIGINAL_MESSAGE)
BasicAcknowledgeablePubsubMessage message, @Headers Map<String, Object> headerMap) {
    log.info("About to go async");
    asyncProcessor(tryMessage, message);
    log.info("Async processor is working");
}

private void asyncProcessor(tryClass tryMessage, BasicAcknowledgeablePubsubMessage message) {
    CompletableFuture.supplyAsync(() -> {
        log.info("Process in Async Fashion");
        // Do heavy lifting here that can take 20ish seconds to consume tryMessage
    });
}

目前处理速度极慢,日志Process in Async Fashion每隔数秒才出现一次,间隔最长可达30秒,从未出现连续输出的情况。

核心问题:导致消费缓慢的原因是什么?

额外疑问:

  • 处理器仅调用asyncProcessor,Future可能运行在其他Executor上,不受上述20×4=80的配置限制,最终是否会引发内存溢出?或者该80的限制会在80个asyncProcessor完成前生效?
  • 是否需要配置spring.cloud.gcp.pubsub.subscription.[subscription-name].flow-control.max-outstanding-element-count=<count>?

回答

一、消费缓慢的核心原因

  1. 默认异步线程池并发度严重不足
    CompletableFuture.supplyAsync()默认使用ForkJoinPool.commonPool(),该线程池的大小通常为CPU核心数-1(比如4核机器仅3个线程)。你的单条业务逻辑需耗时20秒,3个线程同时处理的话,每秒仅能处理0.15条,自然会出现日志间隔长、无法连续输出的情况,完全无法发挥订阅端配置的80线程处理能力。

  2. 未及时ACK消息导致Pub/Sub停止推送
    代码中asyncProcessor仅处理业务逻辑,未调用message.ack()确认消息。Pub/Sub会认为这些消息仍在处理中,当未确认消息数达到阈值后,会停止向订阅端推送新消息,直接导致消费停滞,这是当前消费缓慢的关键原因之一。

  3. 并行拉取能力未被有效利用
    parallel-pull-count=4配置了4个并行拉取请求,但由于业务处理并发不足、消息未及时ACK,拉取到的消息无法快速处理,Pub/Sub会减少返回的消息数量甚至空响应,导致并行拉取的配置完全失效。

二、关于内存溢出与配置的疑问

  1. 默认异步线程池不受Pub/Sub配置约束,存在OOM风险
    你配置的executor-threads和parallel-pull-count仅控制Pub/Sub订阅端拉取、分发消息的线程池,与业务异步线程池无关。CompletableFuture默认的commonPool线程数少,但如果订阅端持续拉取消息(若后续ACK及时),会导致大量未完成的Future任务堆积在队列中,占用内存引发OOM。

  2. 必须配置流量控制参数
    spring.cloud.gcp.pubsub.subscription.[subscription-name].flow-control.max-outstanding-element-count用于限制订阅端未处理(未ACK)的消息数量,避免订阅端拉取过多消息堆积在内存中。结合你的场景,必须配置该参数,建议设置为与业务线程池容量匹配的值,防止内存溢出。

解决建议

  1. 自定义业务异步线程池
    创建与订阅端配置匹配的线程池,充分利用并发能力:
@Bean(name = "asyncProcessorThreadPool")
public Executor asyncProcessorThreadPool() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(80);
    executor.setMaxPoolSize(80);
    executor.setQueueCapacity(100);
    executor.setThreadNamePrefix("AsyncProcessor-");
    executor.initialize();
    return executor;
}

修改asyncProcessor指定该线程池,并添加ACK逻辑:

private void asyncProcessor(tryClass tryMessage, BasicAcknowledgeablePubsubMessage message) {
    CompletableFuture.supplyAsync(() -> {
        log.info("Process in Async Fashion");
        // 业务处理逻辑
        message.ack(); // 处理完成后确认消息
    }, asyncProcessorThreadPool());
}
  1. 配置流量控制参数
    添加如下配置(替换为你的订阅名称):
spring.cloud.gcp.pubsub.subscription.your-subscription-name.flow-control.max-outstanding-element-count=200

该值可根据线程池大小+队列容量调整,避免拉取过多未处理消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 02:47:23