INT-4269: Add Stream & Flux Support for Splitter
JIRA: https://jira.spring.io/browse/INT-4269 * Check if `outputChannel` is `ReactiveStreamsSubscribableChannel` to let back-pressure splitting * Build `Flux` or `Iterator` depending in the `reactive` state * Allow to get a size of the `iterator` if it is possible, for example `XPathMessageSplitter` Add tests Fix raw type and unused import * Add JavaDocs to the `AbstractMessageSplitter#obtainSizeIfPossible()` * Add asserts for the `sequenceSize` populataiton in the `XPathMessageSplitter` * Document `Stream` & `Flux` support in the splitter Minor doc polishing
This commit is contained in:
committed by
Gary Russell
parent
a800d9683a
commit
f112ecbb7e
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2017 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -23,12 +23,19 @@ import java.util.HashMap;
|
||||
import java.util.Iterator;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.function.Function;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import org.reactivestreams.Publisher;
|
||||
|
||||
import org.springframework.integration.channel.ReactiveStreamsSubscribableChannel;
|
||||
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
|
||||
import org.springframework.integration.support.AbstractIntegrationMessageBuilder;
|
||||
import org.springframework.integration.util.FunctionIterator;
|
||||
import org.springframework.messaging.Message;
|
||||
|
||||
import reactor.core.publisher.Flux;
|
||||
|
||||
/**
|
||||
* Base class for Message-splitting handlers.
|
||||
*
|
||||
@@ -57,32 +64,75 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
|
||||
return null;
|
||||
}
|
||||
|
||||
Iterator<Object> iterator;
|
||||
boolean reactive = getOutputChannel() instanceof ReactiveStreamsSubscribableChannel;
|
||||
setAsync(reactive);
|
||||
|
||||
Iterator<Object> iterator = null;
|
||||
Flux<Object> flux = null;
|
||||
|
||||
final int sequenceSize;
|
||||
if (result instanceof Collection) {
|
||||
Collection<Object> items = (Collection<Object>) result;
|
||||
sequenceSize = items.size();
|
||||
iterator = items.iterator();
|
||||
|
||||
if (result instanceof Iterable<?>) {
|
||||
Iterable<Object> iterable = (Iterable<Object>) result;
|
||||
sequenceSize = obtainSizeIfPossible(iterable);
|
||||
if (reactive) {
|
||||
flux = Flux.fromIterable(iterable);
|
||||
}
|
||||
else {
|
||||
iterator = iterable.iterator();
|
||||
}
|
||||
}
|
||||
else if (result.getClass().isArray()) {
|
||||
Object[] items = (Object[]) result;
|
||||
sequenceSize = items.length;
|
||||
iterator = Arrays.asList(items).iterator();
|
||||
}
|
||||
else if (result instanceof Iterable<?>) {
|
||||
sequenceSize = 0;
|
||||
iterator = ((Iterable<Object>) result).iterator();
|
||||
if (reactive) {
|
||||
flux = Flux.fromArray(items);
|
||||
}
|
||||
else {
|
||||
iterator = Arrays.asList(items).iterator();
|
||||
}
|
||||
}
|
||||
else if (result instanceof Iterator<?>) {
|
||||
Iterator<Object> iter = (Iterator<Object>) result;
|
||||
sequenceSize = obtainSizeIfPossible(iter);
|
||||
if (reactive) {
|
||||
flux = Flux.fromIterable(() -> iter);
|
||||
}
|
||||
else {
|
||||
iterator = iter;
|
||||
}
|
||||
}
|
||||
else if (result instanceof Stream<?>) {
|
||||
Stream<Object> stream = ((Stream<Object>) result);
|
||||
sequenceSize = 0;
|
||||
iterator = (Iterator<Object>) result;
|
||||
if (reactive) {
|
||||
flux = Flux.fromStream(stream);
|
||||
}
|
||||
else {
|
||||
iterator = stream.iterator();
|
||||
}
|
||||
}
|
||||
else if (result instanceof Publisher<?>) {
|
||||
Publisher<Object> publisher = (Publisher<Object>) result;
|
||||
sequenceSize = 0;
|
||||
if (reactive) {
|
||||
flux = Flux.from(publisher);
|
||||
}
|
||||
else {
|
||||
iterator = Flux.from((Publisher<Object>) result).toIterable().iterator();
|
||||
}
|
||||
}
|
||||
else {
|
||||
sequenceSize = 1;
|
||||
iterator = Collections.singleton(result).iterator();
|
||||
if (reactive) {
|
||||
flux = Flux.just(result);
|
||||
}
|
||||
else {
|
||||
iterator = Collections.singleton(result).iterator();
|
||||
}
|
||||
}
|
||||
|
||||
if (!iterator.hasNext()) {
|
||||
if (iterator != null && !iterator.hasNext()) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -91,13 +141,42 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
|
||||
messageHeaders = new HashMap<>(messageHeaders);
|
||||
addHeaders(message, messageHeaders);
|
||||
}
|
||||
|
||||
final Map<String, Object> headers = messageHeaders;
|
||||
final Object correlationId = message.getHeaders().getId();
|
||||
final AtomicInteger sequenceNumber = new AtomicInteger(1);
|
||||
|
||||
return new FunctionIterator<Object, AbstractIntegrationMessageBuilder<?>>(iterator,
|
||||
object ->
|
||||
createBuilder(object, headers, correlationId, sequenceNumber.getAndIncrement(), sequenceSize));
|
||||
Function<Object, AbstractIntegrationMessageBuilder<?>> messageBuilderFunction =
|
||||
object -> createBuilder(object, headers, correlationId, sequenceNumber.getAndIncrement(), sequenceSize);
|
||||
|
||||
if (reactive) {
|
||||
return flux.map(messageBuilderFunction);
|
||||
}
|
||||
else {
|
||||
return new FunctionIterator<>(iterator, messageBuilderFunction);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Obtain a size of the provided {@link Iterable}. Default implementation returns
|
||||
* {@link Collection#size()} if the iterable is a collection, or {@code 0} otherwise.
|
||||
* @param iterable the {@link Iterable} to obtain the size
|
||||
* @return the size of the {@link Iterable}
|
||||
* @since 5.0
|
||||
*/
|
||||
protected int obtainSizeIfPossible(Iterable<?> iterable) {
|
||||
return iterable instanceof Collection ? ((Collection<?>) iterable).size() : 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* Obtain a size of the provided {@link Iterator}.
|
||||
* Default implementation returns {@code 0}.
|
||||
* @param iterator the {@link Iterator} to obtain the size
|
||||
* @return the size of the {@link Iterator}
|
||||
* @since 5.0
|
||||
*/
|
||||
protected int obtainSizeIfPossible(Iterator<?> iterator) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
private AbstractIntegrationMessageBuilder<?> createBuilder(Object item, Map<String, Object> headers,
|
||||
@@ -146,10 +225,15 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
|
||||
|
||||
@Override
|
||||
protected void produceOutput(Object result, Message<?> requestMessage) {
|
||||
Iterator<?> iterator = (Iterator<?>) result;
|
||||
while (iterator.hasNext()) {
|
||||
super.produceOutput(iterator.next(), requestMessage);
|
||||
if (result instanceof Iterator<?>) {
|
||||
Iterator<?> iterator = (Iterator<?>) result;
|
||||
while (iterator.hasNext()) {
|
||||
super.produceOutput(iterator.next(), requestMessage);
|
||||
|
||||
}
|
||||
}
|
||||
else {
|
||||
super.produceOutput(result, requestMessage);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2014-2016 the original author or authors.
|
||||
* Copyright 2014-2017 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -26,7 +26,7 @@ import java.util.function.Function;
|
||||
* @author Artem Bilan
|
||||
* @since 4.1
|
||||
*/
|
||||
public final class FunctionIterator<T, V> implements Iterator<V> {
|
||||
public class FunctionIterator<T, V> implements Iterator<V> {
|
||||
|
||||
private final Iterator<T> iterator;
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2013 the original author or authors.
|
||||
* Copyright 2002-2017 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -16,29 +16,40 @@
|
||||
|
||||
package org.springframework.integration.splitter;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.hamcrest.Matchers.is;
|
||||
import static org.hamcrest.Matchers.nullValue;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.assertThat;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.reactivestreams.Subscriber;
|
||||
|
||||
import org.springframework.integration.IntegrationMessageHeaderAccessor;
|
||||
import org.springframework.integration.channel.DirectChannel;
|
||||
import org.springframework.integration.channel.FluxMessageChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.endpoint.EventDrivenConsumer;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.test.StepVerifier;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
* @author Iwein Fuld
|
||||
* @author Gunnar Hillert
|
||||
* @author Artem Bilan
|
||||
*/
|
||||
public class DefaultSplitterTests {
|
||||
|
||||
@@ -65,7 +76,7 @@ public class DefaultSplitterTests {
|
||||
|
||||
@Test
|
||||
public void splitMessageWithCollectionPayload() throws Exception {
|
||||
List<String> payload = Arrays.asList(new String[] { "x", "y", "z" });
|
||||
List<String> payload = Arrays.asList("x", "y", "z");
|
||||
Message<List<String>> message = MessageBuilder.withPayload(payload).build();
|
||||
QueueChannel replyChannel = new QueueChannel();
|
||||
DefaultMessageSplitter splitter = new DefaultMessageSplitter();
|
||||
@@ -108,4 +119,112 @@ public class DefaultSplitterTests {
|
||||
Message<?> output = replyChannel.receive(15);
|
||||
assertThat(output, is(nullValue()));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void splitStream() {
|
||||
Message<?> message = new GenericMessage<>(
|
||||
Stream.generate(Math::random)
|
||||
.limit(10));
|
||||
QueueChannel outputChannel = new QueueChannel();
|
||||
DefaultMessageSplitter splitter = new DefaultMessageSplitter();
|
||||
splitter.setOutputChannel(outputChannel);
|
||||
splitter.handleMessage(message);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
Message<?> reply = outputChannel.receive(0);
|
||||
assertEquals(message.getHeaders().getId(), new IntegrationMessageHeaderAccessor(reply).getCorrelationId());
|
||||
}
|
||||
assertNull(outputChannel.receive(10));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void splitFlux() {
|
||||
Message<?> message = new GenericMessage<>(
|
||||
Flux
|
||||
.generate(() -> 0,
|
||||
(state, sink) -> {
|
||||
if (state == 10) {
|
||||
sink.complete();
|
||||
}
|
||||
else {
|
||||
sink.next(state);
|
||||
}
|
||||
return ++state;
|
||||
}));
|
||||
QueueChannel outputChannel = new QueueChannel();
|
||||
DefaultMessageSplitter splitter = new DefaultMessageSplitter();
|
||||
splitter.setOutputChannel(outputChannel);
|
||||
splitter.handleMessage(message);
|
||||
for (int i = 0; i < 10; i++) {
|
||||
Message<?> reply = outputChannel.receive(0);
|
||||
assertEquals(message.getHeaders().getId(), new IntegrationMessageHeaderAccessor(reply).getCorrelationId());
|
||||
}
|
||||
assertNull(outputChannel.receive(10));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void splitArrayPayloadReactive() {
|
||||
Message<?> message = new GenericMessage<>(new String[] { "x", "y", "z" });
|
||||
FluxMessageChannel replyChannel = new FluxMessageChannel();
|
||||
DefaultMessageSplitter splitter = new DefaultMessageSplitter();
|
||||
splitter.setOutputChannel(replyChannel);
|
||||
|
||||
splitter.handleMessage(message);
|
||||
|
||||
Flux<String> testFlux =
|
||||
Flux.from(replyChannel)
|
||||
.map(Message::getPayload)
|
||||
.cast(String.class);
|
||||
|
||||
StepVerifier.create(testFlux)
|
||||
.expectNext("x", "y", "z")
|
||||
.then(() ->
|
||||
((Subscriber<?>) TestUtils.getPropertyValue(replyChannel, "subscribers", List.class).get(0))
|
||||
.onComplete())
|
||||
.verifyComplete();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void splitStreamReactive() {
|
||||
Message<?> message = new GenericMessage<>(Stream.of("x", "y", "z"));
|
||||
FluxMessageChannel replyChannel = new FluxMessageChannel();
|
||||
DefaultMessageSplitter splitter = new DefaultMessageSplitter();
|
||||
splitter.setOutputChannel(replyChannel);
|
||||
|
||||
splitter.handleMessage(message);
|
||||
|
||||
Flux<String> testFlux =
|
||||
Flux.from(replyChannel)
|
||||
.map(Message::getPayload)
|
||||
.cast(String.class);
|
||||
|
||||
StepVerifier.create(testFlux)
|
||||
.expectNext("x", "y", "z")
|
||||
.then(() ->
|
||||
((Subscriber<?>) TestUtils.getPropertyValue(replyChannel, "subscribers", List.class).get(0))
|
||||
.onComplete())
|
||||
.verifyComplete();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void splitFluxReactive() {
|
||||
Message<?> message = new GenericMessage<>(Flux.just("x", "y", "z"));
|
||||
FluxMessageChannel replyChannel = new FluxMessageChannel();
|
||||
DefaultMessageSplitter splitter = new DefaultMessageSplitter();
|
||||
splitter.setOutputChannel(replyChannel);
|
||||
|
||||
splitter.handleMessage(message);
|
||||
|
||||
Flux<String> testFlux =
|
||||
Flux.from(replyChannel)
|
||||
.map(Message::getPayload)
|
||||
.cast(String.class);
|
||||
|
||||
StepVerifier.create(testFlux)
|
||||
.expectNext("x", "y", "z")
|
||||
.then(() ->
|
||||
((Subscriber<?>) TestUtils.getPropertyValue(replyChannel, "subscribers", List.class).get(0))
|
||||
.onComplete())
|
||||
.verifyComplete();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user