From a2db2787b91dc29c975d967df725f20ecc7bee28 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 19 Oct 2021 17:03:28 +0200 Subject: [PATCH] GH-2235 Fix partitioning issue for collection and add tests --- .../ImplicitFunctionBindingTests.java | 33 +++++++++++++++---- 1 file changed, 26 insertions(+), 7 deletions(-) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 496d7abaf..bba972ea7 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -43,6 +43,8 @@ import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.cloud.function.context.FunctionRegistration; import org.springframework.cloud.function.context.FunctionType; +import org.springframework.cloud.function.context.catalog.FunctionAroundWrapper; +import org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry.FunctionInvocationWrapper; import org.springframework.cloud.function.context.config.ContextFunctionCatalogAutoConfiguration; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.StreamListener; @@ -767,13 +769,13 @@ public class ImplicitFunctionBindingTests { } @Test - public void partitionOnOutputPayloadAsListTest() { + public void partitionOnOutputPayloadAsListAndFunctionAroundWrapperTest() { System.clearProperty("spring.cloud.function.definition"); try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration .getCompleteConfiguration(PojoFunctionConfiguration.class)) - .web(WebApplicationType.NONE).run("--spring.cloud.function.definition=persons", - "--spring.cloud.stream.bindings.persons-out-0.producer.partitionKeyExpression=payload.id", - "--spring.cloud.stream.bindings.persons-out-0.producer.partitionCount=5", + .web(WebApplicationType.NONE).run("--spring.cloud.function.definition=uppercase|persons", + "--spring.cloud.stream.bindings.uppercase|persons-out-0.producer.partitionKeyExpression=payload.id", + "--spring.cloud.stream.bindings.uppercase|persons-out-0.producer.partitionCount=5", "--spring.jmx.enabled=false")) { InputDestination inputDestination = context.getBean(InputDestination.class); @@ -781,9 +783,9 @@ public class ImplicitFunctionBindingTests { Message inputMessage = MessageBuilder.withPayload("Jim Lahey".getBytes()).build(); - inputDestination.send(inputMessage, "persons-in-0"); + inputDestination.send(inputMessage, "uppercasepersons-in-0"); - assertThat(outputDestination.receive(100, "persons-out-0").getHeaders().get("scst_partition")).isEqualTo(3); + assertThat(outputDestination.receive(100, "uppercasepersons-out-0").getHeaders().get("scst_partition")).isEqualTo(3); assertThat(outputDestination.receive(100)).isNull(); } @@ -957,7 +959,7 @@ public class ImplicitFunctionBindingTests { "--spring.cloud.function.definition=supplier")) { OutputDestination outputDestination = context.getBean(OutputDestination.class); - Message result = outputDestination.receive(5000); + Message result = outputDestination.receive(1000, "supplier-out-0"); assertThat(new String(result.getPayload())).isEqualTo("[{\"name\":\"Ricky\",\"id\":1},{\"name\":\"Julien\",\"id\":2}]"); } } @@ -1488,6 +1490,18 @@ public class ImplicitFunctionBindingTests { return x -> x; } + @Bean + public FunctionAroundWrapper faw() { + return new FunctionAroundWrapper() { + + @Override + protected Object doApply(Message input, + FunctionInvocationWrapper targetFunction) { + return targetFunction.apply(input); + } + }; + } + @Bean public Supplier personSupplier() { Person p = new Person(); @@ -1516,6 +1530,11 @@ public class ImplicitFunctionBindingTests { }; } + @Bean + public Function uppercase() { + return v -> v.toUpperCase(); + } + @Bean public Function>> persons() { return x -> {