INT-4740: FileSplitter - Add Headers

JIRA: https://jira.spring.io/browse/INT-3740

For `File` and `String` payloads add `FileHeaders.ORIGINAL_FILE` and `FileHeaders.FILENAME` headers.

Reworked to create headers once only.

Fix typos and Java > 6  API usage
This commit is contained in:
Gary Russell
2015-06-15 09:57:45 -04:00
committed by Artem Bilan
parent 6436aa6d32
commit 3d5f7db4b2
3 changed files with 65 additions and 4 deletions

View File

@@ -19,7 +19,9 @@ package org.springframework.integration.splitter;
import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.Iterator;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
@@ -27,7 +29,6 @@ import org.springframework.integration.support.AbstractIntegrationMessageBuilder
import org.springframework.integration.util.Function;
import org.springframework.integration.util.FunctionIterator;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHeaders;
/**
* Base class for Message-splitting handlers.
@@ -86,21 +87,28 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
return null;
}
final MessageHeaders headers = message.getHeaders();
final Object correlationId = headers.getId();
Map<String, Object> messageHeaders = message.getHeaders();
if (willAddHeaders(message)) {
messageHeaders = new HashMap<String, Object>(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,
new Function<Object, AbstractIntegrationMessageBuilder<?>>() {
@Override
public AbstractIntegrationMessageBuilder<?> apply(Object object) {
return createBuilder(object, headers, correlationId, sequenceNumber.getAndIncrement(),
sequenceSize);
}
});
}
private AbstractIntegrationMessageBuilder<?> createBuilder(Object item, MessageHeaders headers,
private AbstractIntegrationMessageBuilder<?> createBuilder(Object item, Map<String, Object> headers,
Object correlationId, int sequenceNumber, int sequenceSize) {
AbstractIntegrationMessageBuilder<?> builder;
if (item instanceof Message) {
@@ -116,6 +124,26 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
return builder;
}
/**
* Return true if the subclass needs to add headers in the resulting splits.
* If true, {@link #addHeaders} will be called.
* @param message the message.
* @return true
*/
protected boolean willAddHeaders(Message<?> message) {
return false;
}
/**
* Allows subclasses to add extra headers to the output messages. Headers may not be
* removed by this method.
*
* @param message the inbound message.
* @param headers the headers to add messages to.
*/
protected void addHeaders(Message<?> message, Map<String, Object> headers) {
}
@Override
protected boolean shouldCopyRequestHeaders() {
return false;

View File

@@ -30,8 +30,10 @@ import java.nio.charset.Charset;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.NoSuchElementException;
import org.springframework.integration.file.FileHeaders;
import org.springframework.integration.file.splitter.FileSplitter.FileMarker.Mark;
import org.springframework.integration.splitter.AbstractMessageSplitter;
import org.springframework.messaging.Message;
@@ -240,6 +242,32 @@ public class FileSplitter extends AbstractMessageSplitter {
}
}
@Override
protected boolean willAddHeaders(Message<?> message) {
Object payload = message.getPayload();
return payload instanceof File || payload instanceof String;
}
@Override
protected void addHeaders(Message<?> message, Map<String, Object> headers) {
File file = null;
if (message.getPayload() instanceof File) {
file = (File) message.getPayload();
}
else if (message.getPayload() instanceof String) {
file = new File((String) message.getPayload());
}
if (file != null) {
if (!headers.containsKey(FileHeaders.ORIGINAL_FILE)) {
headers.put(FileHeaders.ORIGINAL_FILE, file);
}
if (!headers.containsKey(FileHeaders.FILENAME)) {
headers.put(FileHeaders.FILENAME, file.getName());
}
}
}
public static class FileMarker implements Serializable {
private static final long serialVersionUID = 8514605438145748406L;

View File

@@ -49,6 +49,7 @@ import org.springframework.integration.IntegrationMessageHeaderAccessor;
import org.springframework.integration.annotation.Splitter;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.config.EnableIntegration;
import org.springframework.integration.file.FileHeaders;
import org.springframework.integration.file.splitter.FileSplitter.FileMarker;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
@@ -109,6 +110,8 @@ public class FileSplitterTests {
receive = this.output.receive(10000);
assertNotNull(receive); //äöüß
assertEquals("äöüß", receive.getPayload());
assertEquals(file, receive.getHeaders().get(FileHeaders.ORIGINAL_FILE));
assertEquals(file.getName(), receive.getHeaders().get(FileHeaders.FILENAME));
assertNull(this.output.receive(1));
this.input1.send(new GenericMessage<String>(file.getAbsolutePath()));
@@ -117,6 +120,8 @@ public class FileSplitterTests {
assertEquals(2, receive.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE));
receive = this.output.receive(10000);
assertNotNull(receive); //äöüß
assertEquals(file, receive.getHeaders().get(FileHeaders.ORIGINAL_FILE));
assertEquals(file.getName(), receive.getHeaders().get(FileHeaders.FILENAME));
assertNull(this.output.receive(1));
this.input1.send(new GenericMessage<Reader>(new FileReader(file)));