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
相关产品推荐
相关产品推荐

