diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java b/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java index afe877b51b..f11f0deba6 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/FileWritingMessageHandler.java @@ -23,6 +23,7 @@ import java.io.File; import java.io.FileInputStream; import java.io.FileOutputStream; import java.io.IOException; +import java.io.InputStream; import java.io.OutputStreamWriter; import java.nio.charset.Charset; @@ -81,11 +82,11 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand private static final String LINE_SEPARATOR = System.getProperty("line.separator"); - private volatile String temporaryFileSuffix =".writing"; + private volatile String temporaryFileSuffix = ".writing"; private volatile boolean temporaryFileSuffixSet = false; - private volatile FileExistsMode fileExistsMode = FileExistsMode.REPLACE; + private volatile FileExistsMode fileExistsMode = FileExistsMode.REPLACE; private final Log logger = LogFactory.getLog(this.getClass()); @@ -164,11 +165,11 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand * case the destination exists. For example {@link FileExistsMode#APPEND} * instructs this handler to append data to the existing file rather then * creating a new file for each {@link Message}. - * + *
* If set to {@link FileExistsMode#APPEND}, the adapter will also * create a real instance of the {@link LockRegistry} to ensure that there * is no collisions when multiple threads are writing to the same file. - * + *
* Otherwise the LockRegistry is set to {@link PassThruLockRegistry} which
* has no effect.
*
@@ -179,7 +180,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
Assert.notNull(fileExistsMode, "'fileExistsMode' must not be null.");
this.fileExistsMode = fileExistsMode;
- if (FileExistsMode.APPEND.equals(fileExistsMode)){
+ if (FileExistsMode.APPEND.equals(fileExistsMode)) {
this.lockRegistry = this.lockRegistry instanceof PassThruLockRegistry
? new DefaultLockRegistry()
: this.lockRegistry;
@@ -198,6 +199,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
/**
* If 'true' will append a new-line after each write. It is 'false' by default.
+ *
* @param appendNewLine true if a new-line should be written to the file after payload is written
* @since 4.0.7
*/
@@ -266,7 +268,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
if (!destinationDirectory.exists() && autoCreateDirectory) {
Assert.isTrue(destinationDirectory.mkdirs(),
- "Destination directory [" + destinationDirectory + "] could not be created.");
+ "Destination directory [" + destinationDirectory + "] could not be created.");
}
Assert.isTrue(destinationDirectory.exists(),
@@ -286,11 +288,11 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
Object payload = requestMessage.getPayload();
Assert.notNull(payload, "message payload must not be null");
String generatedFileName = this.fileNameGenerator.generateFileName(requestMessage);
- File originalFileFromHeader = this.retrieveOriginalFileFromHeader(requestMessage);
+ File originalFileFromHeader = retrieveOriginalFileFromHeader(requestMessage);
final File destinationDirectoryToUse = evaluateDestinationDirectoryExpression(requestMessage);
- File tempFile = new File(destinationDirectoryToUse, generatedFileName + temporaryFileSuffix);
+ File tempFile = new File(destinationDirectoryToUse, generatedFileName + this.temporaryFileSuffix);
File resultFile = new File(destinationDirectoryToUse, generatedFileName);
if (FileExistsMode.FAIL.equals(this.fileExistsMode) && resultFile.exists()) {
@@ -306,7 +308,11 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
try {
if (payload instanceof File) {
- resultFile = this.handleFileMessage((File) payload, tempFile, resultFile);
+ resultFile = handleFileMessage((File) payload, tempFile, resultFile);
+ }
+ else if (payload instanceof InputStream) {
+ resultFile = handleInputStreamMessage((InputStream) payload, originalFileFromHeader, tempFile,
+ resultFile);
}
else if (payload instanceof byte[]) {
resultFile = this.handleByteArrayMessage(
@@ -357,17 +363,34 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
}
private File handleFileMessage(final File sourceFile, File tempFile, final File resultFile) throws IOException {
+ if (!FileExistsMode.APPEND.equals(this.fileExistsMode) && this.deleteSourceFiles) {
+ if (sourceFile.renameTo(resultFile)) {
+ return resultFile;
+ }
+ if (logger.isInfoEnabled()) {
+ logger.info(String.format("Failed to move file '%s'. Using copy and delete fallback.",
+ sourceFile.getAbsolutePath()));
+ }
+ }
+ final BufferedInputStream bis = new BufferedInputStream(new FileInputStream(sourceFile));
+ return handleInputStreamMessage(bis, sourceFile, tempFile, resultFile);
+ }
+
+ private File handleInputStreamMessage(final InputStream sourceFileInputStream, File originalFile, File tempFile,
+ final File resultFile) throws IOException {
if (FileExistsMode.APPEND.equals(this.fileExistsMode)) {
File fileToWriteTo = this.determineFileToWrite(resultFile, tempFile);
final BufferedOutputStream bos = new BufferedOutputStream(new FileOutputStream(fileToWriteTo, true));
- final BufferedInputStream bis = new BufferedInputStream(new FileInputStream(sourceFile));
- WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry, fileToWriteTo.getAbsolutePath()){
+
+ WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry,
+ fileToWriteTo.getAbsolutePath()) {
+
@Override
protected void whileLocked() throws IOException {
try {
byte[] buffer = new byte[StreamUtils.BUFFER_SIZE];
int bytesRead = -1;
- while ((bytesRead = bis.read(buffer)) != -1) {
+ while ((bytesRead = sourceFileInputStream.read(buffer)) != -1) {
bos.write(buffer, 0, bytesRead);
}
if (FileWritingMessageHandler.this.appendNewLine) {
@@ -377,7 +400,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
}
finally {
try {
- bis.close();
+ sourceFileInputStream.close();
}
catch (IOException ex) {
}
@@ -388,29 +411,20 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
}
}
}
+
};
whileLockedProcessor.doWhileLocked();
- this.cleanUpAfterCopy(fileToWriteTo, resultFile, sourceFile);
+ cleanUpAfterCopy(fileToWriteTo, resultFile, originalFile);
return resultFile;
}
else {
- if (this.deleteSourceFiles) {
- if (sourceFile.renameTo(resultFile)) {
- return resultFile;
- }
- if (logger.isInfoEnabled()) {
- logger.info(String.format("Failed to move file '%s'. Using copy and delete fallback.",
- sourceFile.getAbsolutePath()));
- }
- }
BufferedOutputStream bos = new BufferedOutputStream(new FileOutputStream(tempFile));
- BufferedInputStream bis = new BufferedInputStream(new FileInputStream(sourceFile));
try {
byte[] buffer = new byte[StreamUtils.BUFFER_SIZE];
int bytesRead = -1;
- while ((bytesRead = bis.read(buffer)) != -1) {
+ while ((bytesRead = sourceFileInputStream.read(buffer)) != -1) {
bos.write(buffer, 0, bytesRead);
}
if (this.appendNewLine) {
@@ -420,7 +434,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
}
finally {
try {
- bis.close();
+ sourceFileInputStream.close();
}
catch (IOException ex) {
}
@@ -430,18 +444,21 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
catch (IOException ex) {
}
}
- this.cleanUpAfterCopy(tempFile, resultFile, sourceFile);
+ cleanUpAfterCopy(tempFile, resultFile, originalFile);
return resultFile;
}
}
- private File handleByteArrayMessage(final byte[] bytes, File originalFile, File tempFile, final File resultFile) throws IOException {
+ private File handleByteArrayMessage(final byte[] bytes, File originalFile, File tempFile, final File resultFile)
+ throws IOException {
File fileToWriteTo = this.determineFileToWrite(resultFile, tempFile);
final boolean append = FileExistsMode.APPEND.equals(this.fileExistsMode);
final BufferedOutputStream bos = new BufferedOutputStream(new FileOutputStream(fileToWriteTo, append));
- WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry, fileToWriteTo.getAbsolutePath()){
+ WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry,
+ fileToWriteTo.getAbsolutePath()) {
+
@Override
protected void whileLocked() throws IOException {
try {
@@ -465,13 +482,17 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
return resultFile;
}
- private File handleStringMessage(final String content, File originalFile, File tempFile, final File resultFile) throws IOException {
+ private File handleStringMessage(final String content, File originalFile, File tempFile, final File resultFile)
+ throws IOException {
File fileToWriteTo = this.determineFileToWrite(resultFile, tempFile);
final boolean append = FileExistsMode.APPEND.equals(this.fileExistsMode);
- final BufferedWriter writer = new BufferedWriter(new OutputStreamWriter(new FileOutputStream(fileToWriteTo, append), this.charset));
- WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry, fileToWriteTo.getAbsolutePath()){
+ final BufferedWriter writer =
+ new BufferedWriter(new OutputStreamWriter(new FileOutputStream(fileToWriteTo, append), this.charset));
+ WhileLockedProcessor whileLockedProcessor = new WhileLockedProcessor(this.lockRegistry,
+ fileToWriteTo.getAbsolutePath()) {
+
@Override
protected void whileLocked() throws IOException {
try {
@@ -497,7 +518,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
return resultFile;
}
- private File determineFileToWrite(File resultFile, File tempFile){
+ private File determineFileToWrite(File resultFile, File tempFile) {
final File fileToWriteTo;
@@ -512,12 +533,12 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
break;
default:
throw new IllegalStateException("Unsupported FileExistsMode "
- + this.fileExistsMode);
+ + this.fileExistsMode);
}
return fileToWriteTo;
}
- private void cleanUpAfterCopy(File fileToWriteTo, File resultFile, File originalFile) throws IOException{
+ private void cleanUpAfterCopy(File fileToWriteTo, File resultFile, File originalFile) throws IOException {
if (!FileExistsMode.APPEND.equals(this.fileExistsMode) && StringUtils.hasText(this.temporaryFileSuffix)) {
this.renameTo(fileToWriteTo, resultFile);
}
@@ -527,24 +548,27 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
}
}
- private void renameTo(File tempFile, File resultFile) throws IOException{
+ private void renameTo(File tempFile, File resultFile) throws IOException {
Assert.notNull(resultFile, "'resultFile' must not be null");
Assert.notNull(tempFile, "'tempFile' must not be null");
if (resultFile.exists()) {
- if (resultFile.setWritable(true, false) && resultFile.delete()){
+ if (resultFile.setWritable(true, false) && resultFile.delete()) {
if (!tempFile.renameTo(resultFile)) {
- throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() + "' to '" + resultFile.getAbsolutePath() + "'");
+ throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() +
+ "' to '" + resultFile.getAbsolutePath() + "'");
}
}
else {
- throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() + "' to '" + resultFile.getAbsolutePath() +
+ throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() +
+ "' to '" + resultFile.getAbsolutePath() +
"' since '" + resultFile.getName() + "' is not writable or can not be deleted");
}
}
else {
if (!tempFile.renameTo(resultFile)) {
- throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() + "' to '" + resultFile.getAbsolutePath() + "'");
+ throw new IOException("Failed to rename file '" + tempFile.getAbsolutePath() +
+ "' to '" + resultFile.getAbsolutePath() + "'");
}
}
}
@@ -558,7 +582,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
if (destinationDirectoryToUse == null) {
throw new IllegalStateException(String.format("The provided " +
- "destinationDirectoryExpression (%s) must not resolve to null.",
+ "destinationDirectoryExpression (%s) must not resolve to null.",
this.destinationDirectoryExpression.getExpressionString()));
}
else if (destinationDirectoryToUse instanceof String) {
@@ -572,7 +596,8 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand
}
else if (destinationDirectoryToUse instanceof File) {
destinationDirectory = (File) destinationDirectoryToUse;
- } else {
+ }
+ else {
throw new IllegalStateException(String.format("The provided " +
"destinationDirectoryExpression (%s) must be of type " +
"java.io.File or be a String.", this.destinationDirectoryExpression.getExpressionString()));
diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/FileWritingMessageHandlerTests.java b/spring-integration-file/src/test/java/org/springframework/integration/file/FileWritingMessageHandlerTests.java
index 2a697f8241..bcbb49fdcb 100644
--- a/spring-integration-file/src/test/java/org/springframework/integration/file/FileWritingMessageHandlerTests.java
+++ b/spring-integration-file/src/test/java/org/springframework/integration/file/FileWritingMessageHandlerTests.java
@@ -29,8 +29,10 @@ import static org.junit.Assert.fail;
import static org.mockito.Mockito.mock;
import java.io.File;
+import java.io.FileInputStream;
import java.io.FileOutputStream;
import java.io.IOException;
+import java.io.InputStream;
import java.io.UnsupportedEncodingException;
import org.junit.Before;
@@ -38,7 +40,6 @@ import org.junit.Ignore;
import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.TemporaryFolder;
-
import org.springframework.beans.factory.BeanFactory;
import org.springframework.integration.channel.NullChannel;
import org.springframework.integration.channel.QueueChannel;
@@ -55,6 +56,7 @@ import org.springframework.util.FileCopyUtils;
* @author Alex Peters
* @author Gary Russell
* @author Tony Falabella
+ * @author Gunnar Hillert
*/
public class FileWritingMessageHandlerTests {
@@ -171,6 +173,29 @@ public class FileWritingMessageHandlerTests {
assertFileContentIs(result, SAMPLE_CONTENT + System.getProperty("line.separator"));
}
+ @Test
+ public void inputStreamPayloadCopiedToNewFile() throws Exception {
+ InputStream is = new FileInputStream(sourceFile);
+ Message> message = MessageBuilder.withPayload(is).build();
+ QueueChannel output = new QueueChannel();
+ handler.setOutputChannel(output);
+ handler.handleMessage(message);
+ Message> result = output.receive(0);
+ assertFileContentIsMatching(result);
+ }
+
+ @Test
+ public void inputStreamPayloadCopiedToNewFileWithNewLines() throws Exception {
+ InputStream is = new FileInputStream(sourceFile);
+ Message> message = MessageBuilder.withPayload(is).build();
+ QueueChannel output = new QueueChannel();
+ handler.setOutputChannel(output);
+ handler.setAppendNewLine(true);
+ handler.handleMessage(message);
+ Message> result = output.receive(0);
+ assertFileContentIs(result, SAMPLE_CONTENT + System.getProperty("line.separator"));
+ }
+
@Test @Ignore // INT-3289 ignored because it won't fail on all OS
public void testCreateDirFail() {
File dir = new File("/foo");
@@ -274,6 +299,44 @@ public class FileWritingMessageHandlerTests {
assertFalse(sourceFile.exists());
}
+ @Test
+ public void deleteSourceFileWithInputstreamPayloadAndFileInstanceHeader() throws Exception {
+ QueueChannel output = new QueueChannel();
+ handler.setCharset(DEFAULT_ENCODING);
+ handler.setDeleteSourceFiles(true);
+ handler.setOutputChannel(output);
+
+ InputStream is = new FileInputStream(sourceFile);
+
+ Message> message = MessageBuilder.withPayload(is)
+ .setHeader(FileHeaders.ORIGINAL_FILE, sourceFile)
+ .build();
+ assertTrue(sourceFile.exists());
+ handler.handleMessage(message);
+ Message> result = output.receive(0);
+ assertFileContentIsMatching(result);
+ assertFalse(sourceFile.exists());
+ }
+
+ @Test
+ public void deleteSourceFileWithInputstreamPayloadAndFilePathHeader() throws Exception {
+ QueueChannel output = new QueueChannel();
+ handler.setCharset(DEFAULT_ENCODING);
+ handler.setDeleteSourceFiles(true);
+ handler.setOutputChannel(output);
+
+ InputStream is = new FileInputStream(sourceFile);
+
+ Message> message = MessageBuilder.withPayload(is)
+ .setHeader(FileHeaders.ORIGINAL_FILE, sourceFile.getAbsolutePath())
+ .build();
+ assertTrue(sourceFile.exists());
+ handler.handleMessage(message);
+ Message> result = output.receive(0);
+ assertFileContentIsMatching(result);
+ assertFalse(sourceFile.exists());
+ }
+
@Test
public void customFileNameGenerator() throws Exception {
final String anyFilename = "fooBar.test";
diff --git a/src/reference/asciidoc/file.adoc b/src/reference/asciidoc/file.adoc
index f169c12c44..e86558e2f7 100644
--- a/src/reference/asciidoc/file.adoc
+++ b/src/reference/asciidoc/file.adoc
@@ -261,7 +261,12 @@ IMPORTANT: Specifying the `delay`, `end` or `reopen` attributes, forces the use
=== Writing files
To write messages to the file system you can use a http://static.springsource.org/spring-integration/api/org/springframework/integration/file/FileWritingMessageHandler.html[FileWritingMessageHandler].
-This class can deal with _File_, _String_, or _byte array_ payloads.
+This class can deal with the following payload types:
+
+* _File_,
+* _String_
+* _byte array_
+* _InputStream_ (since _version 4.2_)
You can configure the encoding and the charset that will be used in case of a String payload.
diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc
index 0b170c5766..ad6e1d3301 100644
--- a/src/reference/asciidoc/whats-new.adoc
+++ b/src/reference/asciidoc/whats-new.adoc
@@ -48,6 +48,8 @@ The `ignore-hidden` attribute has been introduced for the `