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

Spring Integration Flow读写PubSub与PostgreSQL性能慢并发优化问询

问题描述

从Google PubSub拉取记录并写入PostgreSQL数据库时,使用Spring Integration Flow性能极差,插入1000条记录耗时约15分钟。当前流程未实现并行执行,已配置spring.cloud.gcp.pubsub.subscriber.executor-threads=4和spring.cloud.gcp.pubsub.subscriber.parallel-pull-count=3,但线程仍处于parking状态。

当前读取PubSub Topic的代码:

public IntegrationFlow processEvent() {
        return IntegrationFlow.from(Consumer.class, gateway -> gateway.beanName("onMessage"))
                    .transform(Transformers.fromJson(Alert.class))
                .log(LoggingHandler.Level.DEBUG, "Message Routed to DB store", t -> t.toString())
                .routeToRecipients(r -> r
                        .recipient(persistenceRouteChannel()))                      
                .get();
    }
优化方案

1. 让下游处理通道支持并发

你的PubSub配置是针对拉取阶段的,但如果后续的消息处理通道是单线程的,拉取线程会被阻塞等待处理完成,进而出现parking状态。需要将persistenceRouteChannel配置为并发通道:

@Bean
public MessageChannel persistenceRouteChannel() {
    // 根据服务器CPU核数和数据库承受能力调整线程数
    return new ExecutorChannel(Executors.newFixedThreadPool(8));
}

也可以通过配置文件定义:

spring:
  integration:
    channels:
      persistenceRouteChannel:
        type: executor
        capacity: 1000
        executor:
          core-size: 8
          max-size: 16

2. 用批量插入替代单条写入

单条插入数据库的开销极大,建议通过Spring Integration的聚合组件攒批后批量写入:

public IntegrationFlow processEvent() {
    return IntegrationFlow.from(Consumer.class, gateway -> gateway.beanName("onMessage"))
                .transform(Transformers.fromJson(Alert.class))
                .log(LoggingHandler.Level.DEBUG, "Message Routed to DB store", t -> t.toString())
                // 攒批:满100条或10秒触发一次批量写入
                .aggregate(a -> a
                        .correlationStrategy(m -> "batch")
                        .releaseStrategy(g -> g.size() >= 100 || System.currentTimeMillis() - g.getFirst().getHeaders().getTimestamp() >= 10000)
                        .sendPartialResultOnExpiry(true)
                        .expireGroupsUponCompletion(true))
                .routeToRecipients(r -> r.recipient(persistenceRouteChannel()))                      
                .get();
}

对应的持久化方法接收批量数据:

@Transactional
public void batchSave(List<Alert> alerts) {
    alertRepository.saveAll(alerts);
}

3. 调整PubSub的Ack模式并手动确认

默认自动Ack会在消息接收后立即确认,若下游处理失败会导致消息丢失;同时手动Ack配合并发处理,能避免拉取线程被阻塞:

public IntegrationFlow processEvent() {
    return IntegrationFlow.from(Consumer.class, gateway -> gateway.beanName("onMessage"))
                .transform(pubsubMsg -> {
                    Alert alert = objectMapper.readValue(pubsubMsg.getPayload(), Alert.class);
                    return MessageBuilder.withPayload(alert)
                            .setHeader(GcpPubSubHeaders.ORIGINAL_MESSAGE, pubsubMsg)
                            .build();
                })
                .log(LoggingHandler.Level.DEBUG, "Message Routed to DB store", t -> t.toString())
                .aggregate(...) // 可选攒批步骤
                .handle(message -> {
                    List<Alert> alerts = (List<Alert>) message.getPayload();
                    AckablePubsubMessage originalMsg = message.getHeaders().get(GcpPubSubHeaders.ORIGINAL_MESSAGE, AckablePubsubMessage.class);
                    try {
                        persistenceService.batchSave(alerts);
                        originalMsg.ack();
                    } catch (Exception e) {
                        originalMsg.nack();
                        throw e;
                    }
                })
                .get();
}

同时在配置文件中设置Ack模式:

spring:
  cloud:
    gcp:
      pubsub:
        subscriber:
          ack-mode: MANUAL

4. 扩容数据库连接池

数据库连接池大小会限制并发插入的数量,确保连接池足够支撑并发线程:

spring:
  datasource:
    hikari:
      maximum-pool-size: 16 # 建议和并发线程数匹配或稍大
      minimum-idle: 8

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 11:28:21