From b5dd33cc3c4663d9c26f5a8ccdd8ee3701510635 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 28 Jun 2021 15:53:59 +0200 Subject: [PATCH] GH-2179 provide mechanism to remove short-lived headers Headers such as 'spring.cloud.stream.sendto.destination' are no longer valid once the message has been sent, os it will be removed --- .../stream/function/FunctionConfiguration.java | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index df0e7b378..098c8025a 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -337,6 +337,11 @@ public class FunctionConfiguration { : MessageBuilder.withPayload(value).build(); } + @SuppressWarnings({ "unchecked", "rawtypes" }) + private static Message sanitize(Message inputMessage) { + return MessageBuilder.fromMessage(inputMessage).removeHeader("spring.cloud.stream.sendto.destination").build(); + } + private static class FunctionToDestinationBinder implements InitializingBean, ApplicationContextAware { protected final Log logger = LogFactory.getLog(getClass()); @@ -429,7 +434,12 @@ public class FunctionConfiguration { + inputBindingName + ".consumer.concurrency=" + consumerProperties.getConcurrency() + "'"); } SubscribableChannel inputChannel = this.applicationContext.getBean(inputBindingName, SubscribableChannel.class); - return IntegrationReactiveUtils.messageChannelToFlux(inputChannel); + return IntegrationReactiveUtils.messageChannelToFlux(inputChannel).map(m -> { + if (m instanceof Message) { + m = sanitize(m); + } + return m; + }); }).toArray(Publisher[]::new); Function functionToInvoke = function; @@ -647,6 +657,7 @@ public class FunctionConfiguration { @SuppressWarnings("unchecked") @Override public Object apply(Message message) { + message = sanitize(message); if (message != null && consumerProperties != null) { Map headersMap = (Map) ReflectionUtils .getField(this.headersField, message.getHeaders());