diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSource.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSource.java index eb4543277..727522dd1 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSource.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultPollableMessageSource.java @@ -31,6 +31,7 @@ import org.springframework.core.ParameterizedTypeReference; import org.springframework.integration.StaticMessageHeaderAccessor; import org.springframework.integration.acks.AckUtils; import org.springframework.integration.acks.AcknowledgmentCallback; +import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.core.MessageSource; import org.springframework.integration.core.MessagingTemplate; import org.springframework.integration.support.DefaultErrorMessageStrategy; @@ -62,6 +63,12 @@ import org.springframework.util.Assert; */ public class DefaultPollableMessageSource implements PollableMessageSource, Lifecycle, RetryListener { + private static final DirectChannel dummyChannel = new DirectChannel(); + + static { + dummyChannel.setBeanName("dummy.required.by.nonnull.api"); + } + protected static final ThreadLocal attributesHolder = new ThreadLocal(); private final List interceptors = new ArrayList<>(); @@ -103,7 +110,7 @@ public class DefaultPollableMessageSource implements PollableMessageSource, Life if (result instanceof Message) { Message received = (Message) result; for (ChannelInterceptor interceptor : this.interceptors) { - received = interceptor.preSend(received, null); + received = interceptor.preSend(received, dummyChannel); if (received == null) { return null; }