diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/jsr/partition/PartitionCollectorAdapter.java b/spring-batch-core/src/main/java/org/springframework/batch/core/jsr/partition/PartitionCollectorAdapter.java index 1677a721a..1c6965fdc 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/jsr/partition/PartitionCollectorAdapter.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/jsr/partition/PartitionCollectorAdapter.java @@ -64,5 +64,12 @@ public class PartitionCollectorAdapter implements ChunkListener { @Override public void afterChunkError(ChunkContext context) { + try { + synchronized (partitionQueue) { + partitionQueue.add(collector.collectPartitionData()); + } + } catch (Throwable e) { + throw new BatchRuntimeException("An error occured while collecting data from the PartionCollector", e); + } } } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/jsr/partition/PartitionCollectorAdapterTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/jsr/partition/PartitionCollectorAdapterTests.java index 1dfe1b6e5..9a06d17bb 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/jsr/partition/PartitionCollectorAdapterTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/jsr/partition/PartitionCollectorAdapterTests.java @@ -44,7 +44,7 @@ public class PartitionCollectorAdapterTests { }); adapter.afterChunk(null); - adapter.afterChunk(null); + adapter.afterChunkError(null); adapter.afterChunk(null); assertEquals(3, dataQueue.size());