diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java index fc585a1b..baba04c4 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/adapter/MessagingMessageListenerAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2021 the original author or authors. + * Copyright 2002-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -33,6 +33,7 @@ import org.springframework.amqp.support.converter.MessagingMessageConverter; import org.springframework.core.MethodParameter; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.MessagingException; import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.Headers; @@ -379,7 +380,8 @@ public class MessagingMessageListenerAdapter extends AbstractAdaptableMessageLis * We ignore parameters with type Message because they are not involved with conversion. */ boolean isHeaderOrHeaders = methodParameter.hasParameterAnnotation(Header.class) - || methodParameter.hasParameterAnnotation(Headers.class); + || methodParameter.hasParameterAnnotation(Headers.class) + || methodParameter.getParameterType().equals(MessageHeaders.class); boolean isPayload = methodParameter.hasParameterAnnotation(Payload.class); if (isHeaderOrHeaders && isPayload && MessagingMessageListenerAdapter.this.logger.isWarnEnabled()) { MessagingMessageListenerAdapter.this.logger.warn(this.method.getName() diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/annotation/EnableRabbitIntegrationTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/annotation/EnableRabbitIntegrationTests.java index 850cf110..3b2bee26 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/annotation/EnableRabbitIntegrationTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/annotation/EnableRabbitIntegrationTests.java @@ -120,6 +120,7 @@ import org.springframework.core.convert.support.DefaultConversionService; import org.springframework.core.task.TaskExecutor; import org.springframework.data.web.JsonPath; import org.springframework.lang.NonNull; +import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.GenericMessageConverter; import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.Payload; @@ -560,8 +561,9 @@ public class EnableRabbitIntegrationTests { public void testDifferentTypes() throws InterruptedException { Foo1 foo = new Foo1(); foo.setBar("bar"); + this.service.foos.clear(); this.jsonRabbitTemplate.convertAndSend("differentTypes", foo); - assertThat(this.service.latch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(this.service.dtLatch1.await(10, TimeUnit.SECONDS)).isTrue(); assertThat(this.service.foos.get(0)).isInstanceOf(Foo2.class); assertThat(((Foo2) this.service.foos.get(0)).getBar()).isEqualTo("bar"); assertThat(TestUtils.getPropertyValue(this.registry.getListenerContainer("different"), "concurrentConsumers")).isEqualTo(2); @@ -571,8 +573,9 @@ public class EnableRabbitIntegrationTests { public void testDifferentTypesWithConcurrency() throws InterruptedException { Foo1 foo = new Foo1(); foo.setBar("bar"); - this.jsonRabbitTemplate.convertAndSend("differentTypes", foo); - assertThat(this.service.latch.await(10, TimeUnit.SECONDS)).isTrue(); + this.service.foos.clear(); + this.jsonRabbitTemplate.convertAndSend("differentTypes2", foo); + assertThat(this.service.dtLatch2.await(10, TimeUnit.SECONDS)).isTrue(); assertThat(this.service.foos.get(0)).isInstanceOf(Foo2.class); assertThat(((Foo2) this.service.foos.get(0)).getBar()).isEqualTo("bar"); MessageListenerContainer container = this.registry.getListenerContainer("differentWithConcurrency"); @@ -584,8 +587,9 @@ public class EnableRabbitIntegrationTests { public void testDifferentTypesWithVariableConcurrency() throws InterruptedException { Foo1 foo = new Foo1(); foo.setBar("bar"); - this.jsonRabbitTemplate.convertAndSend("differentTypes", foo); - assertThat(this.service.latch.await(10, TimeUnit.SECONDS)).isTrue(); + this.service.foos.clear(); + this.jsonRabbitTemplate.convertAndSend("differentTypes3", foo); + assertThat(this.service.dtLatch3.await(10, TimeUnit.SECONDS)).isTrue(); assertThat(this.service.foos.get(0)).isInstanceOf(Foo2.class); assertThat(((Foo2) this.service.foos.get(0)).getBar()).isEqualTo("bar"); MessageListenerContainer container = this.registry.getListenerContainer("differentWithVariableConcurrency"); @@ -1086,7 +1090,11 @@ public class EnableRabbitIntegrationTests { final List foos = new ArrayList<>(); - final CountDownLatch latch = new CountDownLatch(1); + final CountDownLatch dtLatch1 = new CountDownLatch(1); + + final CountDownLatch dtLatch2 = new CountDownLatch(1); + + final CountDownLatch dtLatch3 = new CountDownLatch(1); final CountDownLatch validationLatch = new CountDownLatch(1); @@ -1237,21 +1245,21 @@ public class EnableRabbitIntegrationTests { containerFactory = "jsonListenerContainerFactoryNoClassMapper") public void handleDifferent(@Validated Foo2 foo) { foos.add(foo); - latch.countDown(); + dtLatch1.countDown(); } @RabbitListener(id = "differentWithConcurrency", queues = "differentTypes2", - containerFactory = "jsonListenerContainerFactory", concurrency = "#{3}") - public void handleDifferentWithConcurrency(Foo2 foo) { + containerFactory = "jsonListenerContainerFactoryNoClassMapper", concurrency = "#{3}") + public void handleDifferentWithConcurrency(Foo2 foo, MessageHeaders headers) { foos.add(foo); - latch.countDown(); + dtLatch2.countDown(); } @RabbitListener(id = "differentWithVariableConcurrency", queues = "differentTypes3", containerFactory = "jsonListenerContainerFactory", concurrency = "3-4") public void handleDifferentWithVariableConcurrency(Foo2 foo) { foos.add(foo); - latch.countDown(); + dtLatch3.countDown(); } @RabbitListener(id = "notStarted", containerFactory = "rabbitAutoStartFalseListenerContainerFactory",