From 5bf152915ddb9b751027b644eadef2fd202f277c Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 12 Jul 2023 18:39:22 -0400 Subject: [PATCH] GH-2770: sendto header and key extraction When sendto header is used for dynamic destinations and a partition key extractor is given for binder based partitioning, then the partition key extractor is not invoked when publishing the message. Addressing this issue. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2770 --- .../ImplicitFunctionBindingTests.java | 54 +++++++++++++++++++ .../function/FunctionConfiguration.java | 7 +++ 2 files changed, 61 insertions(+) diff --git a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index ddae8e574..d59d8bb9f 100644 --- a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -59,11 +59,13 @@ import org.springframework.cloud.stream.binder.test.OutputDestination; import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; import org.springframework.cloud.stream.binding.BindingsLifecycleController; import org.springframework.cloud.stream.binding.BindingsLifecycleController.State; +import org.springframework.cloud.stream.binding.DefaultPartitioningInterceptor; import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel; import org.springframework.context.ApplicationListener; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.core.ResolvableType; +import org.springframework.integration.channel.AbstractMessageChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.scheduling.PollerMetadata; @@ -73,6 +75,7 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.SubscribableChannel; +import org.springframework.messaging.support.ChannelInterceptor; import org.springframework.messaging.support.GenericMessage; import org.springframework.scheduling.support.PeriodicTrigger; @@ -962,6 +965,49 @@ public class ImplicitFunctionBindingTests { } } + // See this issue for more context on this test: https://github.com/spring-cloud/spring-cloud-stream/issues/2770 + @Test + void testSendToDestinationWhenPartitionsEnabled() { + System.clearProperty("spring.cloud.function.definition"); + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(ImperativeSendToDestinationConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false", + // partition-key-extractor also works below, but for easy testing we picked partition-key-expression since + // partition-key-extractor requires defining a bean. + "--spring.cloud.stream.bindings.aa.producer.partition-key-expression=payload")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context.getBean(OutputDestination.class); + Message inputMessage = MessageBuilder.withPayload("aa".getBytes()).build(); + inputDestination.send(inputMessage, "echo-in-0"); + Message receivedMessage = outputDestination.receive(1000, "aa"); + + StreamBridge streamBridge = context.getBean(StreamBridge.class); + MessageChannel aa = streamBridge.resolveDestination("aa", null, null); + List interceptors1 = ((AbstractMessageChannel) aa).getInterceptors(); + + assertThat(interceptors1.size()).isEqualTo(1); + assertThat(interceptors1.get(0)).isInstanceOf(DefaultPartitioningInterceptor.class); + + assertThat(receivedMessage.getPayload()).isEqualTo("aa".getBytes()); + assertThat(receivedMessage.getHeaders().get("spring.cloud.stream.sendto.destination")).isNotNull(); + + inputMessage = MessageBuilder.withPayload("bb".getBytes()).build(); + inputDestination.send(inputMessage, "echo-in-0"); + receivedMessage = outputDestination.receive(1000, "bb"); + + MessageChannel bb = streamBridge.resolveDestination("bb", null, null); + List interceptors2 = ((AbstractMessageChannel) aa).getInterceptors(); + + assertThat(interceptors2.size()).isEqualTo(1); + assertThat(interceptors2.get(0)).isInstanceOf(DefaultPartitioningInterceptor.class); + + assertThat(receivedMessage.getPayload()).isEqualTo("bb".getBytes()); + assertThat(receivedMessage.getHeaders().get("spring.cloud.stream.sendto.destination")).isNotNull(); + } + } + + @Test @ExtendWith(OutputCaptureExtension.class) void testReactiveConsumerWithConcurrencyGreaterThanOneLogsWarning(CapturedOutput output) { @@ -1404,6 +1450,14 @@ public class ImplicitFunctionBindingTests { } } + @EnableAutoConfiguration + public static class ImperativeSendToDestinationConfiguration { + @Bean + public Function> echo() { + return s -> MessageBuilder.withPayload(s).setHeader("spring.cloud.stream.sendto.destination", s).build(); + } + } + @EnableAutoConfiguration public static class SCF_GH_409Configuration { diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index a06e46686..0d46a5a32 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -71,6 +71,7 @@ import org.springframework.cloud.stream.binder.PartitionHandler; import org.springframework.cloud.stream.binder.ProducerProperties; import org.springframework.cloud.stream.binder.ProducerProperties.PollerProperties; import org.springframework.cloud.stream.binding.BindableProxyFactory; +import org.springframework.cloud.stream.binding.DefaultPartitioningInterceptor; import org.springframework.cloud.stream.binding.NewDestinationBindingCallback; import org.springframework.cloud.stream.binding.SupportedBindableFeatures; import org.springframework.cloud.stream.config.BinderFactoryAutoConfiguration; @@ -681,6 +682,12 @@ public class FunctionConfiguration { if (result instanceof Message messageResult && messageResult.getHeaders().get("spring.cloud.stream.sendto.destination") != null) { String destinationName = (String) messageResult.getHeaders().get("spring.cloud.stream.sendto.destination"); MessageChannel outputChannel = streamBridge.resolveDestination(destinationName, producerProperties, null); + BindingProperties bindingProperties = serviceProperties.getBindingProperties(destinationName); + ProducerProperties sendToBindingProducerProperties = bindingProperties.getProducer(); + if (sendToBindingProducerProperties != null && sendToBindingProducerProperties.isPartitioned()) { + ((AbstractMessageChannel) outputChannel) + .addInterceptor(new DefaultPartitioningInterceptor(bindingProperties, applicationContext.getBeanFactory())); + } if (logger.isInfoEnabled()) { logger.info("Output message is sent to '" + destinationName + "' destination"); }