diff --git a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/listener/StreamListenerContainer.java b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/listener/StreamListenerContainer.java index 56bcf5f5..dc83799c 100644 --- a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/listener/StreamListenerContainer.java +++ b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/listener/StreamListenerContainer.java @@ -104,6 +104,7 @@ public class StreamListenerContainer implements MessageListenerContainer, BeanNa * @param messageConverter the converter. */ public void setMessageConverter(StreamMessageConverter messageConverter) { + Assert.notNull(messageConverter, "'messageConverter' cannot be null"); this.messageConverter = messageConverter; } diff --git a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/support/converter/DefaultStreamMessageConverter.java b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/support/converter/DefaultStreamMessageConverter.java index 76c6c0d7..9f68c6e2 100644 --- a/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/support/converter/DefaultStreamMessageConverter.java +++ b/spring-rabbit-stream/src/main/java/org/springframework/rabbit/stream/support/converter/DefaultStreamMessageConverter.java @@ -143,20 +143,23 @@ public class DefaultStreamMessageConverter implements StreamMessageConverter { StreamMessageProperties mProps) { Properties properties = streamMessage.getProperties(); - JavaUtils.INSTANCE - .acceptIfNotNull(properties.getMessageIdAsString(), mProps::setMessageId) - .acceptIfNotNull(properties.getUserId(), usr -> mProps.setUserId(new String(usr, this.charset))) - .acceptIfNotNull(properties.getTo(), mProps::setTo) - .acceptIfNotNull(properties.getSubject(), mProps::setSubject) - .acceptIfNotNull(properties.getReplyTo(), mProps::setReplyTo) - .acceptIfNotNull(properties.getCorrelationIdAsString(), mProps::setCorrelationId) - .acceptIfNotNull(properties.getContentType(), mProps::setContentType) - .acceptIfNotNull(properties.getContentEncoding(), mProps::setContentEncoding) - .acceptIfNotNull(properties.getAbsoluteExpiryTime(), exp -> mProps.setExpiration(Long.toString(exp))) - .acceptIfNotNull(properties.getCreationTime(), mProps::setCreationTime) - .acceptIfNotNull(properties.getGroupId(), mProps::setGroupId) - .acceptIfNotNull(properties.getGroupSequence(), mProps::setGroupSequence) - .acceptIfNotNull(properties.getReplyToGroupId(), mProps::setReplyToGroupId); + if (properties != null) { + JavaUtils.INSTANCE + .acceptIfNotNull(properties.getMessageIdAsString(), mProps::setMessageId) + .acceptIfNotNull(properties.getUserId(), usr -> mProps.setUserId(new String(usr, this.charset))) + .acceptIfNotNull(properties.getTo(), mProps::setTo) + .acceptIfNotNull(properties.getSubject(), mProps::setSubject) + .acceptIfNotNull(properties.getReplyTo(), mProps::setReplyTo) + .acceptIfNotNull(properties.getCorrelationIdAsString(), mProps::setCorrelationId) + .acceptIfNotNull(properties.getContentType(), mProps::setContentType) + .acceptIfNotNull(properties.getContentEncoding(), mProps::setContentEncoding) + .acceptIfNotNull(properties.getAbsoluteExpiryTime(), + exp -> mProps.setExpiration(Long.toString(exp))) + .acceptIfNotNull(properties.getCreationTime(), mProps::setCreationTime) + .acceptIfNotNull(properties.getGroupId(), mProps::setGroupId) + .acceptIfNotNull(properties.getGroupSequence(), mProps::setGroupSequence) + .acceptIfNotNull(properties.getReplyToGroupId(), mProps::setReplyToGroupId); + } Map applicationProperties = streamMessage.getApplicationProperties(); if (applicationProperties != null) { mProps.getHeaders().putAll(applicationProperties);