如何创建两个带不同消息转换器的PubSubTemplate Bean?
GCP Pub/Sub多PubSubTemplate Bean被覆盖问题解决
问题描述
需要创建两个PubSubTemplate Bean,分别配置不同的消息转换器:
pubSubTemplateForUserCreation:使用JacksonPubSubMessageConverter处理JSON格式的UserChangeEvent消息pubSubTemplateForAuditTracker:使用SimplePubSubMessageConverter处理String格式的审计消息
配置完成后通过@Qualifier指定注入目标Bean,但容器启动后仅存在一个PubSubTemplate实例,消息转换器为Jackson版本,pubSubTemplateForAuditTracker Bean被覆盖,导致审计订阅者接收String消息时报错。
原配置代码
PubSubTemplateConfig.java
@Configuration public class PubSubTemplateConfig { @Bean public PubSubTemplate pubSubTemplateForUserCreation(PubSubPublisherTemplate pubSubPublisherTemplate, PubSubSubscriberTemplate pubSubSubscriberTemplate) { PubSubTemplate template = new PubSubTemplate(pubSubPublisherTemplate, pubSubSubscriberTemplate); template.setMessageConverter(new JacksonPubSubMessageConverter(getObjectMapper())); return template; } private ObjectMapper getObjectMapper() { ObjectMapper objectMapper = new ObjectMapper(); objectMapper.registerModule(new JavaTimeModule()); return objectMapper; } @Bean public PubSubTemplate pubSubTemplateForAuditTracker(PubSubPublisherTemplate pubSubPublisherTemplate, PubSubSubscriberTemplate pubSubSubscriberTemplate) { PubSubTemplate template = new PubSubTemplate(pubSubPublisherTemplate, pubSubSubscriberTemplate); template.setMessageConverter(new SimplePubSubMessageConverter()); return template; } }
AuditsubscriptioncriptionConfiguration.java
@Configuration public class AuditsubscriptioncriptionConfiguration { @Value("${subscriptioncription.auditsubscriptioncription}") private String subscription; @Bean("pubsubAuditInputChannel") public MessageChannel pubsubAuditInputChannel() { return new DirectChannel(); } @Bean public PubSubInboundChannelAdapter auditMessageChannelAdapter(@Qualifier("pubsubAuditInputChannel") MessageChannel pubsubAuditInputChannel, @Qualifier("pubSubTemplateForAuditTracker") PubSubTemplate pubSubTemplateForAuditTracker) { PubSubInboundChannelAdapter adapter = new PubSubInboundChannelAdapter(pubSubTemplateForAuditTracker, subscription); adapter.setOutputChannel(pubsubAuditInputChannel); adapter.setPayloadType(String.class); adapter.setAckMode(AckMode.MANUAL); return adapter; } }
UserSubscriptionConfiguration.java
@Configuration public class UserSubscriptionConfiguration { @Value("${subscription.userSubscriber}") private String subscriber; @Bean("pubsubInputChannel") public MessageChannel pubsubInputChannel() { return new DirectChannel(); } @Bean public PubSubInboundChannelAdapter userMessageChannelAdapter(@Qualifier("pubsubInputChannel") MessageChannel pubsubInputChannel, @Qualifier("pubSubTemplateForUserCreation") PubSubTemplate pubSubTemplateForUserCreation) { PubSubInboundChannelAdapter adapter = new PubSubInboundChannelAdapter(pubSubTemplateForUserCreation, subscriber); adapter.setOutputChannel(pubsubInputChannel); adapter.setPayloadType(UserChangeEvent.class); adapter.setAckMode(AckMode.MANUAL); return adapter; } }
解决方案
问题核心是两个PubSubTemplate复用了同一个全局的PubSubPublisherTemplate和PubSubSubscriberTemplate实例,这两个底层模板内部持有消息转换器的引用,后设置的转换器会覆盖之前的配置。要实现完全独立的PubSubTemplate,需为每个模板单独创建对应的底层模板实例。
修改后的PubSubTemplateConfig.java:
@Configuration public class PubSubTemplateConfig { @Bean public PubSubTemplate pubSubTemplateForUserCreation(PubSubClientFactory clientFactory) { // 为用户模板单独创建Publisher和Subscriber实例 PubSubPublisherTemplate publisherTemplate = new PubSubPublisherTemplate(clientFactory); PubSubSubscriberTemplate subscriberTemplate = new PubSubSubscriberTemplate(clientFactory); PubSubTemplate template = new PubSubTemplate(publisherTemplate, subscriberTemplate); template.setMessageConverter(new JacksonPubSubMessageConverter(getObjectMapper())); return template; } private ObjectMapper getObjectMapper() { ObjectMapper objectMapper = new ObjectMapper(); objectMapper.registerModule(new JavaTimeModule()); return objectMapper; } @Bean public PubSubTemplate pubSubTemplateForAuditTracker(PubSubClientFactory clientFactory) { // 为审计模板单独创建Publisher和Subscriber实例 PubSubPublisherTemplate publisherTemplate = new PubSubPublisherTemplate(clientFactory); PubSubSubscriberTemplate subscriberTemplate = new PubSubSubscriberTemplate(clientFactory); PubSubTemplate template = new PubSubTemplate(publisherTemplate, subscriberTemplate); template.setMessageConverter(new SimplePubSubMessageConverter()); return template; } }
关键说明
- 不再依赖全局的
PubSubPublisherTemplate和PubSubSubscriberTemplate,而是通过PubSubClientFactory为每个PubSubTemplate单独实例化底层模板,确保每个模板的消息转换器配置完全隔离,不会互相覆盖。 - 原有的订阅者配置类无需修改,
@Qualifier的注入逻辑依然有效,两个订阅者会各自使用对应的PubSubTemplate处理消息。
内容的提问来源于stack exchange,提问作者Tech_sharma
相关产品推荐
相关产品推荐

