如何修复向GCP Pub/Sub发送多消息时的通道未正确关闭错误
解决Spring Boot集成GCP Pub/Sub时的Managed Channel未关闭问题
问题根源
这个错误的核心是GCP Pub/Sub的Publisher实例绑定的gRPC ManagedChannel未被正确关闭,主要来自两个问题:
- 响应式流处理疏漏:Controller中调用
publishMessageToGcpTopic后未订阅返回的Mono,导致消息发布的异步操作脱离响应式上下文,可能在资源未释放时就被终止。 - 未主动回收Publisher实例:
DefaultPublisherFactory会缓存创建的Publisher,应用关闭或不再使用时,若未主动关闭这些实例,对应的gRPC Channel会出现泄漏。
修复方案
1. 调整配置类,暴露DefaultPublisherFactory实例
修改GCPPublisherConfig,让publisherFactory方法返回具体的DefaultPublisherFactory类型,方便后续获取缓存的Publisher实例:
@Configuration @RequiredArgsConstructor public class GCPPublisherConfig { private final AppProperties appProperties; @Bean public DefaultPublisherFactory publisherFactory(CredentialsProvider defaultCredentialsProvider) { DefaultPublisherFactory factory = new DefaultPublisherFactory(() -> appProperties.getPubsubProjectId()); factory.setEnableMessageOrdering(true); factory.setCredentialsProvider(defaultCredentialsProvider); return factory; } @Bean public PubSubTemplate pubSubTemplate( DefaultPublisherFactory publisherFactory, SubscriberFactory subscriberFactory) { return new PubSubTemplate(publisherFactory, subscriberFactory); } }
2. 添加Shutdown Handler,主动关闭所有Publisher
创建一个Bean,在应用关闭时关闭DefaultPublisherFactory中缓存的所有Publisher实例:
@Component @RequiredArgsConstructor public class PubSubShutdownHandler { private final DefaultPublisherFactory publisherFactory; @PreDestroy public void shutdownPubSubPublishers() { publisherFactory.getPublishers().forEach(publisher -> { try { publisher.shutdown(); if (!publisher.awaitTermination(10, TimeUnit.SECONDS)) { publisher.shutdownNow(); } } catch (InterruptedException e) { publisher.shutdownNow(); Thread.currentThread().interrupt(); } }); } }
3. 修正Controller中的响应式流处理
原来的Controller直接调用publishMessageToGcpTopic但未订阅返回的Mono,导致异步操作未被纳入响应式上下文。需要将其改为flatMap的一部分,确保消息发布完成后再继续流的处理:
@GetMapping("/publish/messages") public Mono<Boolean> testMessages(){ return appService.getData() .flatMap(it -> { EdslMessage edslMessage = EdslMessage.builder().build(); // 将publish的Mono纳入流中,确保异步操作完成 return messagePublisher.publishMessageToGcpTopic(edslMessage, String.valueOf(it)) .map(result -> true); }) .switchIfEmpty(Mono.defer(() -> { log.warn("No entries found for date less than current date"); return Mono.just(false); })) .reduce((a, b) -> a && b); }
补充说明
DefaultPublisherFactory的getPublishers()方法会返回所有已创建的Publisher实例,调用它们的shutdown()和awaitTermination()方法可以确保gRPC Channel被正确关闭。- 响应式流中必须订阅所有异步操作的
Mono/Flux,否则这些操作可能会被提前终止,导致资源泄漏或消息丢失。
内容的提问来源于stack exchange,提问作者RKodakandla
相关产品推荐
相关产品推荐

