【问题标题】:validate raw message against schema for method annotated with jmslistener针对使用 jmslistener 注释的方法的模式验证原始消息
【发布时间】:2018-04-11 22:28:58
【问题描述】:

我需要对所有 jms 侦听器应用一些预检查和常见步骤,例如根据模式(JSON 模式)验证原始消息。示例 -

@Component
public class MyService {

    @JmsListener(destination = "myDestination")
    public void processOrder(Order order) { ... }
}

现在,在 spring 将 Message 从队列转换为 Order 之前,我需要执行以下操作 -

  1. 将带有标头的原始消息记录到自定义记录器中。
  2. 根据 json 模式验证 json 消息(文本消息)(为了简单起见,假设我这里只有一个模式)
  3. 如果架构验证失败,记录错误并抛出异常
  4. 如果 schema 验证通过,继续控制 spring 进行转换并继续处理 order 方法。

spring JMS 架构是否提供任何方式来注入上述需求? 我知道 AOP 会出现,但我不确定它是否可以与 @JmsListener 一起使用。

【问题讨论】:

    标签: spring spring-boot jms spring-jms spring-io


    【解决方案1】:

    一种相当简单的技术是在侦听器容器工厂上将autoStartup 设置为false

    然后,使用JmsListenerEndpointRegistry bean 获取侦听器容器。

    然后getMessageListener(),将其包装在AOP代理和setMessageListener()中。

    然后启动容器。

    可能有更优雅的方式,但我认为您必须深入了解侦听器创建代码的内容,这非常复杂。

    编辑

    Spring Boot 示例:

    @SpringBootApplication
    public class So49682934Application {
    
        private final Logger logger = LoggerFactory.getLogger(getClass());
    
        public static void main(String[] args) {
            SpringApplication.run(So49682934Application.class, args);
        }
    
        @JmsListener(id = "listener1", destination = "so49682934")
        public void listen(Foo foo) {
            logger.info(foo.toString());
        }
    
        @Bean
        public ApplicationRunner runner(JmsListenerEndpointRegistry registry, JmsTemplate template) {
            return args -> {
                DefaultMessageListenerContainer container =
                        (DefaultMessageListenerContainer) registry.getListenerContainer("listener1");
                Object listener = container.getMessageListener();
                ProxyFactory pf = new ProxyFactory(listener);
                NameMatchMethodPointcutAdvisor advisor = new NameMatchMethodPointcutAdvisor(new MyJmsInterceptor());
                advisor.addMethodName("onMessage");
                pf.addAdvisor(advisor);
                container.setMessageListener(pf.getProxy());
                registry.start();
                Thread.sleep(5_000);
                Foo foo = new Foo("baz");
                template.convertAndSend("so49682934", foo);
            };
        }
    
        @Bean
        public MessageConverter converter() {
            MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter();
            converter.setTargetType(MessageType.TEXT);
            converter.setTypeIdPropertyName("typeId");
            return converter;
        }
    
        public static class MyJmsInterceptor implements MethodInterceptor {
    
            private final Logger logger = LoggerFactory.getLogger(getClass());
    
            @Override
            public Object invoke(MethodInvocation invocation) throws Throwable {
                Message message = (Message) invocation.getArguments()[0];
                logger.info(message.toString());
                // validate
                return invocation.proceed();
            }
    
        }
    
        public static class Foo {
    
            private String bar;
    
            public Foo() {
                super();
            }
    
            public Foo(String bar) {
                this.bar = bar;
            }
    
            public String getBar() {
                return this.bar;
            }
    
            public void setBar(String bar) {
                this.bar = bar;
            }
    
            @Override
            public String toString() {
                return "Foo [bar=" + this.bar + "]";
            }
    
        }
    
    }
    

    spring.jms.listener.auto-startup=false
    

    m2018-04-06 11:42:04.859 INFO 59745 --- [enerContainer-1] e.So49682934Application$MyJmsInterceptor : ActiveMQTextMessage {commandId = 5, responseRequired = true, messageId = ID:gollum.local-60138-1523029319662 -4:2:1:1:1, originalDestination = null, originalTransactionId = null, producerId = ID:gollum.local-60138-1523029319662-4:2:1:1, destination = queue://so49682934, transactionId = null ,过期 = 0,时间戳 = 1523029324849,到达 = 0,brokerInTime = 1523029324849,brokerOutTime = 1523029324853,correlationId = null,replyTo = null,persistent = true,type = null,priority = 4,groupID = null,groupSequence = 0,targetConsumerId = null,compressed = false,userID = null,content = null,marshalledProperties = null,dataStructure = null,redeliveryCounter = 0,size = 1050,properties = {typeId=com.example.So49682934Application$Foo},readOnlyProperties = true,readOnlyBody = true, droppable = false, jmsXGroupFirstForConsumer = false, text = {"bar":"baz"}}

    2018-04-06 11:42:04.882 INFO 59745 --- [enerContainer-1] ication$$EnhancerBySpringCGLIB$$e29327b8 : Foo [bar=baz]

    EDIT2

    这里是如何通过基础设施来做到这一点...

    @SpringBootApplication
    @EnableJms
    public class So496829341Application {
    
        private final Logger logger = LoggerFactory.getLogger(getClass());
    
        public static void main(String[] args) {
            SpringApplication.run(So496829341Application.class, args);
        }
    
        @JmsListener(id = "listen1", destination="so496829341")
        public void listen(Foo foo) {
            logger.info(foo.toString());
        }
    
        @Bean
        public ApplicationRunner runner(JmsTemplate template) {
            return args -> {
                Thread.sleep(5_000);
                template.convertAndSend("so496829341", new Foo("baz"));
            };
        }
    
        @Bean
        public MessageConverter converter() {
            MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter();
            converter.setTargetType(MessageType.TEXT);
            converter.setTypeIdPropertyName("typeId");
            return converter;
        }
    
        @Bean(JmsListenerConfigUtils.JMS_LISTENER_ANNOTATION_PROCESSOR_BEAN_NAME)
        public static JmsListenerAnnotationBeanPostProcessor bpp() {
            return new JmsListenerAnnotationBeanPostProcessor() {
    
                @Override
                protected MethodJmsListenerEndpoint createMethodJmsListenerEndpoint() {
                    return new MethodJmsListenerEndpoint() {
    
                        @Override
                        protected MessagingMessageListenerAdapter createMessageListener(
                                MessageListenerContainer container) {
                            MessagingMessageListenerAdapter listener = super.createMessageListener(container);
                            ProxyFactory pf = new ProxyFactory(listener);
                            pf.setProxyTargetClass(true);
                            NameMatchMethodPointcutAdvisor advisor = new NameMatchMethodPointcutAdvisor(new MyJmsInterceptor());
                            advisor.addMethodName("onMessage");
                            pf.addAdvisor(advisor);
                            return (MessagingMessageListenerAdapter) pf.getProxy();
                        }
    
                    };
                }
    
            };
        }
    
        public static class MyJmsInterceptor implements MethodInterceptor {
    
            private final Logger logger = LoggerFactory.getLogger(getClass());
    
            @Override
            public Object invoke(MethodInvocation invocation) throws Throwable {
                Message message = (Message) invocation.getArguments()[0];
                logger.info(message.toString());
                // validate
                return invocation.proceed();
            }
    
        }
    
        public static class Foo {
    
            private String bar;
    
            public Foo() {
                super();
            }
    
            public Foo(String bar) {
                this.bar = bar;
            }
    
            public String getBar() {
                return this.bar;
            }
    
            public void setBar(String bar) {
                this.bar = bar;
            }
    
            @Override
            public String toString() {
                return "Foo [bar=" + this.bar + "]";
            }
    
        }
    
    }
    

    注意:BPP 必须是静态的,并且@EnableJms 是必需的,因为存在此 BPP 会禁用引导。

    2018-04-06 13:44:41.607 INFO 82669 --- [enerContainer-1] .So496829341Application$MyJmsInterceptor : ActiveMQTextMessage {commandId = 5, responseRequired = true, messageId = ID:gollum.local-63685-1523036676402- 4:2:1:1:1, originalDestination = null, originalTransactionId = null, producerId = ID:gollum.local-63685-1523036676402-4:2:1:1, destination = queue://so496829341, transactionId = null,过期 = 0,时间戳 = 1523036681598,到达 = 0,brokerInTime = 1523036681598,brokerOutTime = 1523036681602,correlationId = null,replyTo = null,persistent = true,type = null,priority = 4,groupID = null,groupSequence = 0,targetConsumerId = null,compressed = false,userID = null,content = null,marshalledProperties = null,dataStructure = null,redeliveryCounter = 0,size = 1050,properties = {typeId=com.example.So496829341Application$Foo},readOnlyProperties = true,readOnlyBody =真,droppable = 假,jmsXGroupFirstForConsumer = 假,文本 = {"bar":"baz"}}

    2018-04-06 13:44:41.634 INFO 82669 --- [enerContainer-1] ication$$EnhancerBySpringCGLIB$$$9ff4b13f : Foo [bar=baz]

    EDIT3

    避免 AOP...

    @SpringBootApplication
    @EnableJms
    public class So496829341Application {
    
        private final Logger logger = LoggerFactory.getLogger(getClass());
    
        public static void main(String[] args) {
            SpringApplication.run(So496829341Application.class, args);
        }
    
        @JmsListener(id = "listen1", destination="so496829341")
        public void listen(Foo foo) {
            logger.info(foo.toString());
        }
    
        @Bean
        public ApplicationRunner runner(JmsTemplate template) {
            return args -> {
                Thread.sleep(5_000);
                template.convertAndSend("so496829341", new Foo("baz"));
            };
        }
    
        @Bean
        public MessageConverter converter() {
            MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter();
            converter.setTargetType(MessageType.TEXT);
            converter.setTypeIdPropertyName("typeId");
            return converter;
        }
    
        @Bean(JmsListenerConfigUtils.JMS_LISTENER_ANNOTATION_PROCESSOR_BEAN_NAME)
        public static JmsListenerAnnotationBeanPostProcessor bpp() {
            return new JmsListenerAnnotationBeanPostProcessor() {
    
                @Override
                protected MethodJmsListenerEndpoint createMethodJmsListenerEndpoint() {
                    return new MethodJmsListenerEndpoint() {
    
                        @Override
                        protected MessagingMessageListenerAdapter createMessageListener(
                                MessageListenerContainer container) {
                            final MessagingMessageListenerAdapter listener = super.createMessageListener(container);
                            return new MessagingMessageListenerAdapter() {
    
                                @Override
                                public void onMessage(Message jmsMessage, Session session) throws JMSException {
                                    logger.info(jmsMessage.toString());
                                    // validate
                                    listener.onMessage(jmsMessage, session);
                                }
    
                            };
                        }
    
                    };
                }
    
            };
        }
    
        public static class Foo {
    
            private String bar;
    
            public Foo() {
                super();
            }
    
            public Foo(String bar) {
                this.bar = bar;
            }
    
            public String getBar() {
                return this.bar;
            }
    
            public void setBar(String bar) {
                this.bar = bar;
            }
    
            @Override
            public String toString() {
                return "Foo [bar=" + this.bar + "]";
            }
    
        }
    
    }
    

    EDIT4

    要访问监听器方法上的其他注解是可以的,但是需要反射才能得到Method的引用...

    @JmsListener(id = "listen1", destination="so496829341")
    @Schema("foo.bar")
    public void listen(Foo foo) {
        logger.info(foo.toString());
    }
    
    @Target({ElementType.METHOD, ElementType.ANNOTATION_TYPE})
    @Retention(RetentionPolicy.RUNTIME)
    @Inherited
    @Documented
    public @interface Schema {
    
        String value();
    
    }
    
    @Bean(JmsListenerConfigUtils.JMS_LISTENER_ANNOTATION_PROCESSOR_BEAN_NAME)
    public static JmsListenerAnnotationBeanPostProcessor bpp() {
        return new JmsListenerAnnotationBeanPostProcessor() {
    
            @Override
            protected MethodJmsListenerEndpoint createMethodJmsListenerEndpoint() {
                return new MethodJmsListenerEndpoint() {
    
                    @Override
                    protected MessagingMessageListenerAdapter createMessageListener(
                            MessageListenerContainer container) {
                        final MessagingMessageListenerAdapter listener = super.createMessageListener(container);
                        InvocableHandlerMethod handlerMethod =
                                (InvocableHandlerMethod) new DirectFieldAccessor(listener)
                                        .getPropertyValue("handlerMethod");
                        final Schema schema = AnnotationUtils.getAnnotation(handlerMethod.getMethod(), Schema.class);
                        return new MessagingMessageListenerAdapter() {
    
                            @Override
                            public void onMessage(Message jmsMessage, Session session) throws JMSException {
                                logger.info(jmsMessage.toString());
                                logger.info(schema.value());
                                // validate
                                listener.onMessage(jmsMessage, session);
                            }
    
                        };
                    }
    
                };
            }
    
        };
    }
    

    【讨论】:

    • 您可以分享任何文档或示例吗?另外,有没有办法可以扩展@JmsListener 的功能?
    • 查看我的编辑以获得简单的解决方案;是的,你可以扩展@JmsListener,但它非常复杂——你必须扩展JmsListenerAnnotationBeanPostProcessor,并实现一个自定义JmsListenerEndpoint(可能是标准MethodJmsListenerEndpoint的一个子类)。
    • 我进行了第二次编辑以展示如何改写@JmsListener bean 后处理器。第三次编辑做同样的事情,但避免 AOP;只需将适配器包装在另一个适配器中即可。
    • 感谢您的代码,我会尝试并告诉您。您能否帮助我提供可以提供这种级别的实施洞察力的文档或教程?
    • 除了reference manual 中的内容之外,它并没有真正记录在案。这是一种相当先进的技术。我只是碰巧熟悉它,因为我在 spring-amqp 和 spring-kafka 中使用类似的功能,我是项目负责人。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-12-14
    • 1970-01-01
    • 2021-05-27
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多