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 new file mode 100644 index 0000000000..b4f26395f9 --- /dev/null +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/FileSplitter.java @@ -0,0 +1,178 @@ +/* + * Copyright 2015 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.file.splitter; + +import java.io.BufferedReader; +import java.io.File; +import java.io.FileInputStream; +import java.io.FileNotFoundException; +import java.io.FileReader; +import java.io.IOException; +import java.io.InputStream; +import java.io.InputStreamReader; +import java.io.Reader; +import java.nio.charset.Charset; +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; + +import org.springframework.integration.splitter.AbstractMessageSplitter; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHandlingException; + +/** + * The {@link AbstractMessageSplitter} implementation to split the {@link File} + * Message payload to lines. + *

+ * With {@code iterator = true} (defaults to {@code true}) this class produces an {@link Iterator} + * to process file lines on demand from {@link Iterator#next}. + * Otherwise a {@link List} of all lines is returned to the to further + * {@link AbstractMessageSplitter#handleRequestMessage} process. + *

+ * Can accept {@link String} as file path, {@link File}, {@link Reader} or {@link InputStream} + * as payload type. + * All other types are ignored and returned to the {@link AbstractMessageSplitter} as is. + * + * @author Artem Bilan + * @author Gary Russell + * @since 4.1.2 + */ +public class FileSplitter extends AbstractMessageSplitter { + + private final boolean iterator; + + private Charset charset; + + /** + * Construct a splitter where the {@link #splitMessage(Message)} method returns + * an iterator and the file is read line-by-line during iteration. + */ + public FileSplitter() { + this(true); + } + + /** + * 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. + * @param iterator true to return an iterator, false to return a list of lines. + */ + public FileSplitter(boolean iterator) { + this.iterator = iterator; + } + + /** + * Set the charset to be used when reading the file, when something other than the default + * charset is required. + * @param charset the charset. + */ + public void setCharset(Charset charset) { + this.charset = charset; + } + + @Override + protected Object splitMessage(final Message message) { + Object payload = message.getPayload(); + + Reader reader = null; + + if (payload instanceof String) { + try { + reader = new FileReader((String) payload); + } + catch (FileNotFoundException e) { + throw new MessageHandlingException(message, "failed to read file [" + payload + "]", e); + } + } + else if (payload instanceof File) { + try { + if (this.charset == null) { + reader = new FileReader((File) payload); + } + else { + reader = new InputStreamReader(new FileInputStream((File) payload), this.charset); + } + } + catch (FileNotFoundException e) { + throw new MessageHandlingException(message, "failed to read file [" + payload + "]", e); + } + } + else if (payload instanceof InputStream) { + if (this.charset == null) { + reader = new InputStreamReader((InputStream) payload); + } + else { + reader = new InputStreamReader((InputStream) payload, this.charset); + } + } + else if (payload instanceof Reader) { + reader = (Reader) payload; + } + else { + return message; + } + + final BufferedReader bufferedReader = new BufferedReader(reader); + Iterator iterator = new Iterator() { + + @Override + public boolean hasNext() { + try { + boolean ready = bufferedReader.ready(); + if (!ready) { + bufferedReader.close(); + } + return ready; + } + catch (IOException e) { + try { + bufferedReader.close(); + } + catch (IOException e1) {} + throw new MessageHandlingException(message, "IOException while iterating", e); + } + } + + @Override + public String next() { + try { + return bufferedReader.readLine(); + } + catch (IOException e) { + try { + bufferedReader.close(); + } + catch (IOException e1) {} + throw new MessageHandlingException(message, "IOException while iterating", e); + } + } + + }; + + if (this.iterator) { + return iterator; + } + else { + List lines = new ArrayList(); + while (iterator.hasNext()) { + lines.add(iterator.next()); + } + return lines; + } + } + +} diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/package-info.java b/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/package-info.java new file mode 100644 index 0000000000..0cc4ef72db --- /dev/null +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/splitter/package-info.java @@ -0,0 +1,5 @@ +/** + * Provides implementations of + * {@link org.springframework.integration.splitter.AbstractMessageSplitter}. + */ +package org.springframework.integration.file.splitter; diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests-context.xml b/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests-context.xml new file mode 100644 index 0000000000..b0fe39ab9e --- /dev/null +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests-context.xml @@ -0,0 +1,14 @@ + + + + + + + + + + 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 new file mode 100644 index 0000000000..5717749f51 --- /dev/null +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java @@ -0,0 +1,203 @@ +/* + * Copyright 2015 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.file.splitter; + +import static org.hamcrest.Matchers.containsString; +import static org.hamcrest.Matchers.instanceOf; +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.fail; + +import java.io.ByteArrayInputStream; +import java.io.File; +import java.io.FileNotFoundException; +import java.io.FileOutputStream; +import java.io.FileReader; +import java.io.IOException; +import java.io.InputStream; +import java.io.Reader; +import java.nio.charset.Charset; +import java.util.Date; + +import org.junit.AfterClass; +import org.junit.BeforeClass; +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.ImportResource; +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.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessageHandler; +import org.springframework.messaging.PollableChannel; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.test.context.support.AnnotationConfigContextLoader; +import org.springframework.util.FileCopyUtils; + +/** + * @author Artem Bilan + * @author Gary Russell + * @since 4.1.2 + */ +@ContextConfiguration(loader = AnnotationConfigContextLoader.class) +@RunWith(SpringJUnit4ClassRunner.class) +@DirtiesContext +public class FileSplitterTests { + + private static File file; + + static final String SAMPLE_CONTENT = "HelloWorld\näöüß"; + + @Autowired + private MessageChannel input1; + + @Autowired + private MessageChannel input2; + + @Autowired + private MessageChannel input3; + + @Autowired + private PollableChannel output; + + @BeforeClass + public static void setup() throws IOException { + file = File.createTempFile("foo", ".txt"); + FileCopyUtils.copy(SAMPLE_CONTENT.getBytes("UTF-8"), + new FileOutputStream(file, false)); + } + + @AfterClass + public static void tearDown() { + file.delete(); + } + + @Test + public void testFileSplitter() throws Exception { + this.input1.send(new GenericMessage(file)); + Message receive = this.output.receive(10000); + assertNotNull(receive); //HelloWorld + assertEquals(2, receive.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE)); + receive = this.output.receive(10000); + assertNotNull(receive); //äöüß + assertNull(this.output.receive(1)); + + this.input1.send(new GenericMessage(file.getAbsolutePath())); + receive = this.output.receive(10000); + assertNotNull(receive); //HelloWorld + assertEquals(2, receive.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE)); + receive = this.output.receive(10000); + assertNotNull(receive); //äöüß + assertNull(this.output.receive(1)); + + this.input1.send(new GenericMessage(new FileReader(file))); + receive = this.output.receive(10000); + assertNotNull(receive); //HelloWorld + assertEquals(2, receive.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE)); + receive = this.output.receive(10000); + assertNotNull(receive); //äöüß + assertNull(this.output.receive(1)); + + this.input2.send(new GenericMessage(file)); + receive = this.output.receive(10000); + assertNotNull(receive); //HelloWorld + assertEquals(0, receive.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE)); + receive = this.output.receive(10000); + assertNotNull(receive); //äöüß + assertNull(this.output.receive(1)); + + this.input2.send(new GenericMessage(new ByteArrayInputStream(SAMPLE_CONTENT.getBytes("UTF-8")))); + receive = this.output.receive(10000); + assertNotNull(receive); //HelloWorld + assertEquals(0, receive.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE)); + receive = this.output.receive(10000); + assertNotNull(receive); //äöüß + assertNull(this.output.receive(1)); + + try { + this.input2.send(new GenericMessage("bar")); + fail("FileNotFoundException expected"); + } + catch (Exception e) { + assertThat(e.getCause(), instanceOf(FileNotFoundException.class)); + assertThat(e.getMessage(), containsString("failed to read file [bar]")); + } + this.input2.send(new GenericMessage(new Date())); + receive = this.output.receive(10000); + assertNotNull(receive); + assertEquals(1, receive.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE)); + assertThat(receive.getPayload(), instanceOf(Date.class)); + assertNull(this.output.receive(1)); + + this.input3.send(new GenericMessage(file)); + receive = this.output.receive(10000); + assertNotNull(receive); //HelloWorld + assertEquals(0, receive.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE)); + receive = this.output.receive(10000); + assertNotNull(receive); //äöüß + assertNull(this.output.receive(1)); + + this.input3.send(new GenericMessage(new ByteArrayInputStream(SAMPLE_CONTENT.getBytes("UTF-8")))); + receive = this.output.receive(10000); + assertNotNull(receive); //HelloWorld + assertEquals(0, receive.getHeaders().get(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE)); + receive = this.output.receive(10000); + assertNotNull(receive); //äöüß + assertNull(this.output.receive(1)); + } + + @Configuration + @EnableIntegration + @ImportResource("classpath:org/springframework/integration/file/splitter/FileSplitterTests-context.xml") + public static class ContextConfiguration { + + @Bean + public PollableChannel output() { + return new QueueChannel(); + } + + @Bean + @Splitter(inputChannel = "input2") + public MessageHandler fileSplitter2() { + FileSplitter fileSplitter = new FileSplitter(true); + fileSplitter.setOutputChannel(output()); + return fileSplitter; + } + + @Bean + @Splitter(inputChannel = "input3") + public MessageHandler fileSplitter3() { + FileSplitter fileSplitter = new FileSplitter(); + fileSplitter.setCharset(Charset.defaultCharset()); + fileSplitter.setOutputChannel(output()); + return fileSplitter; + } + + } + +}