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
This commit is contained in:
Soby Chacko
2023-07-12 18:39:22 -04:00
committed by Oleg Zhurakousky
parent 403d97c057
commit b83b2c6aa4
2 changed files with 61 additions and 0 deletions

View File

@@ -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<byte[]> inputMessage = MessageBuilder.withPayload("aa".getBytes()).build();
inputDestination.send(inputMessage, "echo-in-0");
Message<byte[]> receivedMessage = outputDestination.receive(1000, "aa");
StreamBridge streamBridge = context.getBean(StreamBridge.class);
MessageChannel aa = streamBridge.resolveDestination("aa", null, null);
List<ChannelInterceptor> 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<ChannelInterceptor> 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<String, Message<String>> echo() {
return s -> MessageBuilder.withPayload(s).setHeader("spring.cloud.stream.sendto.destination", s).build();
}
}
@EnableAutoConfiguration
public static class SCF_GH_409Configuration {

View File

@@ -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");
}