From 302a44ddb6cfc5b7d0373f2687cb81fd756863be Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 15 Oct 2021 19:29:14 +0200 Subject: [PATCH] GH-2235 Fix partitioning support when output is a collection of messages --- .../PartitionAwareFunctionWrapper.java | 35 +++++++++++++++---- 1 file changed, 29 insertions(+), 6 deletions(-) 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 6cf0fb41f..e5ad742ad 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 @@ -16,6 +16,9 @@ package org.springframework.cloud.stream.function; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; import java.util.function.Function; import java.util.function.Supplier; @@ -31,6 +34,7 @@ import org.springframework.expression.spel.support.StandardEvaluationContext; import org.springframework.integration.expression.ExpressionUtils; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; +import org.springframework.util.ObjectUtils; /** * This class is effectively a wrapper which is aware of the stream related partition information @@ -45,7 +49,7 @@ class PartitionAwareFunctionWrapper implements Function, Supplie private final Function function; @SuppressWarnings("rawtypes") - private final Function outputMessageEnricher; + private final Function outputMessageEnricher; PartitionAwareFunctionWrapper(Function function, ConfigurableApplicationContext context, ProducerProperties producerProperties) { this.function = function; @@ -55,13 +59,24 @@ class PartitionAwareFunctionWrapper implements Function, Supplie PartitionHandler partitionHandler = new PartitionHandler(evaluationContext, producerProperties, context.getBeanFactory()); this.outputMessageEnricher = output -> { - if (!(output instanceof Message)) { + if (ObjectUtils.isArray(output) && !(output instanceof byte[])) { + output = Arrays.asList(output); + } + if (output instanceof Iterable) { + Iterable elements = (Iterable) output; + List messages = new ArrayList<>(); + for (Object element : elements) { + if (!(element instanceof Message)) { + element = MessageBuilder.withPayload(element).build(); + } + messages.add(toMessageWithPartitionHeader((Message) element, partitionHandler)); + } + return messages; + } + else if (!(output instanceof Message)) { output = MessageBuilder.withPayload(output).build(); } - int partitionId = partitionHandler.determinePartition((Message) output); - return MessageBuilder - .fromMessage((Message) output) - .setHeader(BinderHeaders.PARTITION_HEADER, partitionId).build(); + return toMessageWithPartitionHeader((Message) output, partitionHandler); }; } else { @@ -69,6 +84,14 @@ class PartitionAwareFunctionWrapper implements Function, Supplie } } + @SuppressWarnings({ "unchecked", "rawtypes" }) + private Message toMessageWithPartitionHeader(Message message, PartitionHandler partitionHandler) { + int partitionId = partitionHandler.determinePartition(message); + return MessageBuilder + .fromMessage(message) + .setHeader(BinderHeaders.PARTITION_HEADER, partitionId).build(); + } + @SuppressWarnings("unchecked") @Override public Object apply(Object input) {