check for null producer properties

This commit is contained in:
ncheema
2022-02-03 13:20:37 -08:00
committed by Oleg Zhurakousky
parent 9ec7cf9b43
commit c1abf65935
2 changed files with 34 additions and 1 deletions

View File

@@ -272,7 +272,7 @@ public final class StreamBridge implements SmartInitializingSingleton {
binder = binderFactory.getBinder(binderName, messageChannel.getClass());
}
if (producerProperties.isPartitioned()) {
if (producerProperties != null && producerProperties.isPartitioned()) {
BindingProperties bindingProperties = this.bindingServiceProperties.getBindingProperties(destinationName);
((AbstractMessageChannel) messageChannel)
.addInterceptor(new DefaultPartitioningInterceptor(bindingProperties, this.applicationContext.getBeanFactory()));

View File

@@ -35,6 +35,7 @@ import org.springframework.boot.WebApplicationType;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry.FunctionInvocationWrapper;
import org.springframework.cloud.stream.binder.test.InputDestination;
import org.springframework.cloud.stream.binder.test.OutputDestination;
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
import org.springframework.cloud.stream.binding.BinderAwareChannelResolver.NewDestinationBindingCallback;
@@ -55,12 +56,15 @@ import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.converter.AbstractMessageConverter;
import org.springframework.messaging.converter.MessageConverter;
import org.springframework.messaging.support.ChannelInterceptor;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.util.MimeType;
import org.springframework.util.MimeTypeUtils;
import org.springframework.util.ReflectionUtils;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.fail;
/**
@@ -467,6 +471,35 @@ public class StreamBridgeTests {
assertThat(context.getBean("callbackVerifier", AtomicBoolean.class)).isTrue();
}
}
@EnableAutoConfiguration
public static class DynamicProducerConfig {
@Bean
public Function<Message<String>, Message<String>> uppercase() {
return msg -> MessageBuilder.withPayload(msg.getPayload().toUpperCase())
.setHeader("spring.cloud.stream.sendto.destination", "dynamicTopic").build();
}
}
@Test
public void testDynamicProducerDestination() {
System.clearProperty("spring.cloud.function.definition");
ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(DynamicProducerConfig.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.function.definition=uppercase",
"--spring.cloud.stream.bindings.uppercase-in-0.destination=upper"
);
InputDestination source = context.getBean(InputDestination.class);
source.send(new GenericMessage<>("John Doe".getBytes()), "upper");
OutputDestination target = context.getBean(OutputDestination.class);
Message<byte[]> message = target.receive(5, "dynamicTopic");
assertNotNull(message);
assertEquals(new String(message.getPayload()), "JOHN DOE");
}
@EnableAutoConfiguration
public static class EmptyConfiguration {