From f772e1ce6b6509f39b5c2be0adc495c3c9376b4b Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 8 Dec 2021 11:28:49 +0100 Subject: [PATCH] GH-2249 Fix partition handling in StreamBridge Resolves #2249 --- .../PartitionAwareFunctionWrapper.java | 6 +++- .../stream/function/StreamBridgeTests.java | 29 +++++++++++++++++++ 2 files changed, 34 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/PartitionAwareFunctionWrapper.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/PartitionAwareFunctionWrapper.java index 01cdbff6a..ef3cc2fde 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/PartitionAwareFunctionWrapper.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/PartitionAwareFunctionWrapper.java @@ -81,7 +81,11 @@ class PartitionAwareFunctionWrapper implements Function, Supplie @Override public Object apply(Object input) { this.setEnhancerIfNecessary(); - return this.function.apply(input); + Object result = this.function.apply(input); + if (!((FunctionInvocationWrapper) this.function).isInputTypePublisher()) { + ((FunctionInvocationWrapper) this.function).setEnhancer(null); + } + return result; } @Override diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java index 366e486ca..78a3d77bc 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java @@ -76,6 +76,35 @@ public class StreamBridgeTests { System.clearProperty("spring.cloud.function.definition"); } + /* + * This test must not result in exception stating "Partition key cannot be null" + * See https://github.com/spring-cloud/spring-cloud-stream/issues/2249 for more details + */ + @Test + public void test_2249() throws Exception { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + EmptyConfiguration.class)).web(WebApplicationType.NONE).run( + "--spring.cloud.stream.source=outputA;outputB", + "--spring.cloud.stream.bindings.outputA-out-0.producer.partition-count=3", + "--spring.cloud.stream.bindings.outputA-out-0.producer.partition-key-expression=headers['partitionKey']", + "--spring.cloud.stream.bindings.outputB-out-0.destination=outputB", + "--spring.cloud.stream.bindings.outputB-out-0.producer.partition-count=3", + "--spring.jmx.enabled=false")) { + StreamBridge streamBridge = context.getBean(StreamBridge.class); + streamBridge.send("outputA-out-0", MessageBuilder.withPayload("A").setHeader("partitionKey", "A").build()); + streamBridge.send("outputB", MessageBuilder.withPayload("B").build()); + streamBridge.send("outputA-out-0", MessageBuilder.withPayload("C").setHeader("partitionKey", "C").build()); + streamBridge.send("outputB", MessageBuilder.withPayload("D").build()); + + OutputDestination output = context.getBean(OutputDestination.class); + assertThat(output.receive(1000, "outputA-out-0").getHeaders().containsKey("scst_partition")).isTrue(); + assertThat(output.receive(1000, "outputB").getHeaders().containsKey("scst_partition")).isFalse(); + assertThat(output.receive(1000, "outputA-out-0").getHeaders().containsKey("scst_partition")).isTrue(); + assertThat(output.receive(1000, "outputB").getHeaders().containsKey("scst_partition")).isFalse(); + } + } + @Test public void testWithOutputContentTypeWildCardBindings() throws Exception { try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration