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

如何创建两个带不同消息转换器的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 13:20:20