【问题标题】:Spring Amqp: Mix SimpleRoutingConnectionFactory with @RabbitListenerSpring Amqp:将 SimpleRoutingConnectionFactory 与 @RabbitListener 混合
【发布时间】:2017-08-04 16:46:50
【问题描述】:

我有一个应用程序会监听多个队列,这些队列在不同的虚拟主机上声明。我使用了一个 SimpleRoutingConnectionFactory 来存储一个 connectionFactoryMap,我希望用@RabbitListener 设置我的监听器。

根据 Spring AMQP 文档:

同样从 1.4 版本开始,您可以配置路由连接 SimpleMessageListenerContainer 中的工厂。在这种情况下,列表 队列名称用作查找键。例如,如果您配置 具有 setQueueNames("foo, bar") 的容器,查找键将是 “[foo,bar]”(无空格)。

我使用了@RabbitListener(queues = "some-key")。不幸的是,spring 抱怨“查找键 [null]”。见下文。

18:52:44.528 警告 --- [cTaskExecutor-1] o.s.a.r.l.SimpleMessageListenerContainer :消费者引发异常, 如果连接工厂支持,处理可以重新启动 java.lang.IllegalStateException:无法确定目标 查找键 [null] 的 ConnectionFactory 位于 org.springframework.amqp.rabbit.connection.AbstractRoutingConnectionFactory.determineTargetConnectionFactory(AbstractRoutingConnectionFactory.java:119) 在 org.springframework.amqp.rabbit.connection.AbstractRoutingConnectionFactory.createConnection(AbstractRoutingConnectionFactory.java:97) 在 org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils$1.createConnection(ConnectionFactoryUtils.java:90) 在 org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils.doGetTransactionalResourceHolder(ConnectionFactoryUtils.java:140) 在 org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils.getTransactionalResourceHolder(ConnectionFactoryUtils.java:76) 在 org.springframework.amqp.rabbit.listener.BlockingQueueConsumer.start(BlockingQueueConsumer.java:472) 在 org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer$AsyncMessageProcessingConsumer.run(SimpleMessageListenerContainer.java:1306) 在 java.lang.Thread.run(Thread.java:745)

  1. 我做错了吗?如果 queues 属性用作查找键(用于连接工厂查找),我应该使用什么来指定我想听哪个队列?

  2. 最终,我希望进行程序化/动态侦听器设置。如果我使用“Programmatic Endpoint Registration”,我应该放弃“Annotation-driven listener endpoints”吗?我喜欢“注释驱动的侦听器端点”,因为侦听器可以有多个消息句柄,具有不同的传入数据类型作为参数,这非常干净整洁。如果我使用 Programmatic Endpoint Registration,我将不得不解析 Message 输入变量,并根据消息类型/内容调用我的特定自定义消息处理程序。

编辑: 嗨,加里, 我稍微修改了您的代码 #2,以便它使用 Jackson2JsonMessageConverter 序列化类对象(在 RabbitTemplate bean 中),并使用它将它们反序列化回对象(在 inboundAdapter 中)。我还删除了@RabbitListener,因为在我的情况下,所有侦听器都将在运行时添加。现在 fooBean 可以毫无问题地接收整数、字符串和 TestData 消息!唯一留下的问题是程序不断报告警告:

"[erContainer#0-1] o.s.a.r.l.SimpleMessageListenerContainer : Consumer 引发异常,如果连接工厂支持,处理可以重新开始

java.lang.IllegalStateException: Cannot determine target ConnectionFactory for lookup key [null]"。完整的堆栈跟踪,请参阅底部。

我错过了什么吗?

@SpringBootApplication
public class App2 implements CommandLineRunner {

    public static void main(String[] args) {
        SpringApplication.run(App2.class, args);
    }

    @Autowired
    private IntegrationFlowContext flowContext;

    @Autowired
    private ConnectionFactory routingCf;

    @Autowired
    private RabbitTemplate template;

    @Override
    public void run(String... args) throws Exception {
        // dynamically add a listener for queue qux
        IntegrationFlow flow = IntegrationFlows.from(Amqp.inboundAdapter(this.routingCf, "qux").messageConverter(new Jackson2JsonMessageConverter()))
                .handle(fooBean())
                .get();
        this.flowContext.registration(flow).register();

        // now test it
        SimpleResourceHolder.bind(this.routingCf, "[qux]");
        this.template.convertAndSend("qux", 42);
        this.template.convertAndSend("qux", "fizbuz");
        this.template.convertAndSend("qux", new TestData(1, "test"));
        SimpleResourceHolder.unbind(this.routingCf);
    }

    @Bean
    RabbitTemplate rabbitTemplate() {
        RabbitTemplate template = new RabbitTemplate(routingCf);
        template.setMessageConverter(new Jackson2JsonMessageConverter());
        return template;
    }

    @Bean
    @Primary
    public ConnectionFactory routingCf() {
        SimpleRoutingConnectionFactory rcf = new SimpleRoutingConnectionFactory();
        Map<Object, ConnectionFactory> map = new HashMap<>();
        map.put("[foo,bar]", routedCf());
        map.put("[baz]", routedCf());
        map.put("[qux]", routedCf());
        rcf.setTargetConnectionFactories(map);
        return rcf;
    }

    @Bean
    public ConnectionFactory routedCf() {
        return new CachingConnectionFactory("127.0.0.1");
    }

    @Bean
    public Foo fooBean() {
        return new Foo();
    }

    public static class Foo {

        @ServiceActivator
        public void handleInteger(Integer in) {
            System.out.println("int: " + in);
        }

        @ServiceActivator
        public void handleString(String in) {
            System.out.println("str: " + in);
        }

        @ServiceActivator
        public void handleData(TestData data) {
            System.out.println("TestData: " + data);
        }
    }
}

完整的堆栈跟踪:

2017-03-15 21:43:06.413  INFO 1003 --- [           main] hello.App2                               : Started App2 in 3.003 seconds (JVM running for 3.69)
2017-03-15 21:43:11.415  WARN 1003 --- [erContainer#0-1] o.s.a.r.l.SimpleMessageListenerContainer : Consumer raised exception, processing can restart if the connection factory supports it

java.lang.IllegalStateException: Cannot determine target ConnectionFactory for lookup key [null]
    at org.springframework.amqp.rabbit.connection.AbstractRoutingConnectionFactory.determineTargetConnectionFactory(AbstractRoutingConnectionFactory.java:119) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.connection.AbstractRoutingConnectionFactory.createConnection(AbstractRoutingConnectionFactory.java:97) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.core.RabbitTemplate.doExecute(RabbitTemplate.java:1430) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.core.RabbitTemplate.execute(RabbitTemplate.java:1411) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.core.RabbitTemplate.execute(RabbitTemplate.java:1387) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.core.RabbitAdmin.initialize(RabbitAdmin.java:500) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.core.RabbitAdmin$11.onCreate(RabbitAdmin.java:419) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.connection.CompositeConnectionListener.onCreate(CompositeConnectionListener.java:33) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.connection.CachingConnectionFactory.createConnection(CachingConnectionFactory.java:571) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils$1.createConnection(ConnectionFactoryUtils.java:90) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils.doGetTransactionalResourceHolder(ConnectionFactoryUtils.java:140) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.connection.ConnectionFactoryUtils.getTransactionalResourceHolder(ConnectionFactoryUtils.java:76) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.listener.BlockingQueueConsumer.start(BlockingQueueConsumer.java:505) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer$AsyncMessageProcessingConsumer.run(SimpleMessageListenerContainer.java:1382) ~[spring-rabbit-1.7.1.RELEASE.jar:na]
    at java.lang.Thread.run(Thread.java:745) [na:1.8.0_112]

【问题讨论】:

    标签: rabbitmq spring-amqp


    【解决方案1】:

    请显示您的配置 - 它适用于我...

    @SpringBootApplication
    public class So42784471Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So42784471Application.class, args);
        }
    
        @Bean
        @Primary
        public ConnectionFactory routing() {
            SimpleRoutingConnectionFactory rcf = new SimpleRoutingConnectionFactory();
            Map<Object, ConnectionFactory> map = new HashMap<>();
            map.put("[foo,bar]", routedCf());
            map.put("[baz]", routedCf());
            rcf.setTargetConnectionFactories(map);
            return rcf;
        }
    
        @Bean
        public ConnectionFactory routedCf() {
            return new CachingConnectionFactory("10.0.0.3");
        }
    
        @RabbitListener(queues = { "foo" , "bar" })
        public void foobar(String in) {
            System.out.println(in);
        }
    
        @RabbitListener(queues = "baz")
        public void bazzer(String in) {
            System.out.println(in);
        }
    
    }
    

    关于您的第二个问题,您可以手动构建端点,但它非常复杂。在 Spring Integration @ServiceActivator 中使用类似功能可能更容易。

    我将尽快更新此答案的详细信息。

    编辑

    这是使用 Spring Integration 技术在运行时动态添加多方法侦听器的更新...

    @SpringBootApplication
    public class So42784471Application implements CommandLineRunner {
    
        public static void main(String[] args) {
            SpringApplication.run(So42784471Application.class, args);
        }
    
        @Autowired
        private IntegrationFlowContext flowContext;
    
        @Autowired
        private ConnectionFactory routingCf;
    
        @Autowired
        private RabbitTemplate template;
    
        @Override
        public void run(String... args) throws Exception {
            // dynamically add a listener for queue qux
            IntegrationFlow flow = IntegrationFlows.from(Amqp.inboundAdapter(this.routingCf, "qux"))
                    .handle(fooBean())
                    .get();
            this.flowContext.registration(flow).register();
    
            // now test it
            SimpleResourceHolder.bind(this.routingCf, "[qux]");
            this.template.convertAndSend("qux", 42);
            this.template.convertAndSend("qux", "fizbuz");
            SimpleResourceHolder.unbind(this.routingCf);
        }
    
    
        @Bean
        @Primary
        public ConnectionFactory routingCf() {
            SimpleRoutingConnectionFactory rcf = new SimpleRoutingConnectionFactory();
            Map<Object, ConnectionFactory> map = new HashMap<>();
            map.put("[foo,bar]", routedCf());
            map.put("[baz]", routedCf());
            map.put("[qux]", routedCf());
            rcf.setTargetConnectionFactories(map);
            return rcf;
        }
    
        @Bean
        public ConnectionFactory routedCf() {
            return new CachingConnectionFactory("10.0.0.3");
        }
    
        @RabbitListener(queues = { "foo" , "bar" })
        public void foobar(String in) {
            System.out.println(in);
        }
    
        @RabbitListener(queues = "baz")
        public void bazzer(String in) {
            System.out.println(in);
        }
    
        @Bean
        public Foo fooBean() {
            return new Foo();
        }
    
        public static class Foo {
    
            @ServiceActivator
            public void handleInteger(Integer in) {
                System.out.println("int: " + in);
            }
    
            @ServiceActivator
            public void handleString(String in) {
                System.out.println("str: " + in);
            }
    
        }
    
    }
    

    【讨论】:

    • 您好,格雷,感谢您的快速回复。我希望你昨天回答我的问题。而你做到了!我编译并运行了您的代码 sn-p #1,并得到“java.lang.IllegalStateException:无法确定查找键 [null] 的目标 ConnectionFactory”。任何想法?我使用的是 Spring boot 1.5.2,当然还创建了队列 foo、bar 和 baz。
    • 嗨,Gary,代码 #2 看起来很有希望。我稍微修改了一下。我更新了我的帖子。请看一看。有一个警告一直困扰着我。我想我错过了一些重要的事情。
    • 发生了一些奇怪的事情 - 根据您的堆栈跟踪,消费者正在尝试通过RoutingConnectionFactory 进行连接,但根据this code,目标连接工厂应该已经被提取。看起来容器是在设置队列名称之前启动的,这对我来说毫无意义。我也在使用 Boot 1.5.2。也许您可以使用调试器来查看发生了什么?
    • 嗨@gary-russell,我试过调试器,但还没有任何线索。我在我的帖子中发布了完整的堆栈跟踪。您是否有一点时间通过在您的计算机上运行我的代码来帮助我,看看您是否可以重现该问题?
    • 我刚刚将您的代码复制到我的应用程序中,它工作得很好(没有错误)。你在getConnectionFactory() 中设置了断点吗?查找没有找到目标工厂吗?如果你想将你的完整项目发布到 github 或其他地方,我可以看看。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-07-02
    • 1970-01-01
    • 1970-01-01
    • 2012-07-24
    • 1970-01-01
    • 2018-12-27
    • 1970-01-01
    相关资源
    最近更新 更多