From 51b52ea3a66a04bfffc83fedd6d36435103d1b50 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 5 Jun 2015 18:42:12 -0400 Subject: [PATCH] INT-3727: FileSplitter - add SOF/EOF Capability JIRA: https://jira.spring.io/browse/INT-3727 Add a configuration option such that the `FileSplitter` can emit Start/End messages when splitting. Polishing according PR comments Add toString to FileMarker --- .../integration/util/FunctionIterator.java | 5 +- .../file/splitter/FileSplitter.java | 94 +++++++++++++++++-- .../file/splitter/FileSplitterTests.java | 25 +++++ 3 files changed, 114 insertions(+), 10 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/FunctionIterator.java b/spring-integration-core/src/main/java/org/springframework/integration/util/FunctionIterator.java index b315220e78..c80dde293a 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/util/FunctionIterator.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/FunctionIterator.java @@ -53,10 +53,7 @@ public final class FunctionIterator implements Iterator { @Override public V next() { - if (this.hasNext()) { - return this.function.apply(this.iterator.next()); - } - throw new NoSuchElementException(); + return this.function.apply(this.iterator.next()); } } diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/FileSplitter.java b/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/FileSplitter.java index b4f26395f9..282ac23ae0 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/FileSplitter.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/FileSplitter.java @@ -25,11 +25,13 @@ import java.io.IOException; import java.io.InputStream; import java.io.InputStreamReader; import java.io.Reader; +import java.io.Serializable; import java.nio.charset.Charset; import java.util.ArrayList; import java.util.Iterator; import java.util.List; +import org.springframework.integration.file.splitter.FileSplitter.FileMarker.Mark; import org.springframework.integration.splitter.AbstractMessageSplitter; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandlingException; @@ -55,6 +57,8 @@ public class FileSplitter extends AbstractMessageSplitter { private final boolean iterator; + private final boolean markers; + private Charset charset; /** @@ -62,7 +66,7 @@ public class FileSplitter extends AbstractMessageSplitter { * an iterator and the file is read line-by-line during iteration. */ public FileSplitter() { - this(true); + this(true, false); } /** @@ -72,7 +76,25 @@ public class FileSplitter extends AbstractMessageSplitter { * @param iterator true to return an iterator, false to return a list of lines. */ public FileSplitter(boolean iterator) { + this(iterator, false); + } + + /** + * Construct a splitter where the {@link #splitMessage(Message)} method returns + * an iterator, and the file is read line-by-line during iteration, or a list + * of lines from the file. When file markers are enabled (START/END) + * {@link #setApplySequence(boolean) applySequence} is false by default. If enabled, + * the markers are included in the sequence size. + * @param iterator true to return an iterator, false to return a list of lines. + * @param markers true to emit start of file/end of file marker messages before/after the data. + * @since 1.4.5 + */ + public FileSplitter(boolean iterator, boolean markers) { this.iterator = iterator; + this.markers = markers; + if (markers) { + setApplySequence(false); + } } /** @@ -90,9 +112,12 @@ public class FileSplitter extends AbstractMessageSplitter { Reader reader = null; + final String filePath; + if (payload instanceof String) { try { reader = new FileReader((String) payload); + filePath = (String) payload; } catch (FileNotFoundException e) { throw new MessageHandlingException(message, "failed to read file [" + payload + "]", e); @@ -106,6 +131,7 @@ public class FileSplitter extends AbstractMessageSplitter { else { reader = new InputStreamReader(new FileInputStream((File) payload), this.charset); } + filePath = ((File) payload).getAbsolutePath(); } catch (FileNotFoundException e) { throw new MessageHandlingException(message, "failed to read file [" + payload + "]", e); @@ -118,25 +144,38 @@ public class FileSplitter extends AbstractMessageSplitter { else { reader = new InputStreamReader((InputStream) payload, this.charset); } + filePath = ":stream:"; } else if (payload instanceof Reader) { reader = (Reader) payload; + filePath = ":reader:"; } else { return message; } final BufferedReader bufferedReader = new BufferedReader(reader); - Iterator iterator = new Iterator() { + Iterator iterator = new Iterator() { + + boolean markers = FileSplitter.this.markers; + + boolean sof = markers; + + boolean eof; + + boolean done; @Override public boolean hasNext() { try { - boolean ready = bufferedReader.ready(); + boolean ready = !this.done && bufferedReader.ready(); if (!ready) { + if (this.markers) { + this.eof = true; + } bufferedReader.close(); } - return ready; + return this.sof || ready || this.eof; } catch (IOException e) { try { @@ -148,7 +187,17 @@ public class FileSplitter extends AbstractMessageSplitter { } @Override - public String next() { + public Object next() { + if (this.sof) { + this.sof = false; + return new FileMarker(filePath, Mark.START); + } + if (this.eof) { + this.eof = false; + this.markers = false; + this.done = true; + return new FileMarker(filePath, Mark.END); + } try { return bufferedReader.readLine(); } @@ -167,7 +216,7 @@ public class FileSplitter extends AbstractMessageSplitter { return iterator; } else { - List lines = new ArrayList(); + List lines = new ArrayList(); while (iterator.hasNext()) { lines.add(iterator.next()); } @@ -175,4 +224,37 @@ public class FileSplitter extends AbstractMessageSplitter { } } + public static class FileMarker implements Serializable { + + private static final long serialVersionUID = 8514605438145748406L; + + public enum Mark implements Serializable { + START, + END + } + + private final String filePath; + + private final Mark mark; + + public FileMarker(String filePath, Mark mark) { + this.filePath = filePath; + this.mark = mark; + } + + public String getFilePath() { + return filePath; + } + + public Mark getMark() { + return mark; + } + + @Override + public String toString() { + return "FileMarker [filePath=" + filePath + ", mark=" + mark + "]"; + } + + } + } diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java b/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java index 5717749f51..10bb34d2ed 100644 --- a/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java @@ -35,6 +35,7 @@ import java.io.Reader; import java.nio.charset.Charset; import java.util.Date; +import org.hamcrest.Matchers; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.Test; @@ -48,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.splitter.FileSplitter.FileMarker; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; @@ -171,6 +173,29 @@ public class FileSplitterTests { assertNull(this.output.receive(1)); } + @Test + public void testMarkers() { + QueueChannel outputChannel = new QueueChannel(); + FileSplitter splitter = new FileSplitter(true, true); + splitter.setOutputChannel(outputChannel); + splitter.handleMessage(new GenericMessage(file)); + Message received = outputChannel.receive(0); + assertNotNull(received); + assertNull(received.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE)); + assertThat(received.getPayload(), Matchers.instanceOf(FileSplitter.FileMarker.class)); + FileMarker fileMarker = (FileSplitter.FileMarker) received.getPayload(); + assertEquals(FileSplitter.FileMarker.Mark.START, fileMarker.getMark()); + assertEquals(file.getAbsolutePath(), fileMarker.getFilePath()); + assertNotNull(outputChannel.receive(0)); + assertNotNull(outputChannel.receive(0)); + received = outputChannel.receive(0); + assertNotNull(received); + assertThat(received.getPayload(), Matchers.instanceOf(FileSplitter.FileMarker.class)); + fileMarker = (FileSplitter.FileMarker) received.getPayload(); + assertEquals(FileSplitter.FileMarker.Mark.END, fileMarker.getMark()); + assertEquals(file.getAbsolutePath(), fileMarker.getFilePath()); + } + @Configuration @EnableIntegration @ImportResource("classpath:org/springframework/integration/file/splitter/FileSplitterTests-context.xml")