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 f8e2a416a0..7ef15b4309 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 @@ -391,7 +391,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand } @Override - public void stop() { + public synchronized void stop() { if (this.flushTask != null) { this.flushTask.cancel(true); this.flushTask = null; @@ -945,7 +945,7 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand this.lock = lock; } - private void close() { + private boolean close() { try { this.lock.lockInterruptibly(); try { @@ -959,9 +959,11 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand catch (IOException e) { // ignore } + return true; } catch (InterruptedException e1) { Thread.currentThread().interrupt(); + return false; } finally { this.lock.unlock(); @@ -982,10 +984,14 @@ public class FileWritingMessageHandler extends AbstractReplyProducingMessageHand FileState state = entry.getValue(); if (state.lastWrite < expired || (!FileWritingMessageHandler.this.flushWhenIdle && state.firstWrite < expired)) { - iterator.remove(); - state.close(); - if (FileWritingMessageHandler.this.logger.isDebugEnabled()) { - FileWritingMessageHandler.this.logger.debug("Flushed: " + entry.getKey()); + if (state.close()) { + if (FileWritingMessageHandler.this.logger.isDebugEnabled()) { + FileWritingMessageHandler.this.logger.debug("Flushed: " + entry.getKey()); + } + iterator.remove(); + } + else { + break; // interrupted } } }