diff --git a/spring-batch-integration/pom.xml b/spring-batch-integration/pom.xml index 95ce94291..e9679d12a 100644 --- a/spring-batch-integration/pom.xml +++ b/spring-batch-integration/pom.xml @@ -52,6 +52,13 @@ ${project.version} compile + + org.springframework.batch + org.springframework.batch.test + + ${project.version} + test + javax.jms com.springsource.javax.jms diff --git a/spring-batch-integration/src/main/java/org/springframework/batch/integration/chunk/ChunkMessageChannelItemWriter.java b/spring-batch-integration/src/main/java/org/springframework/batch/integration/chunk/ChunkMessageChannelItemWriter.java index 5326e0ec3..bcc73f0fb 100644 --- a/spring-batch-integration/src/main/java/org/springframework/batch/integration/chunk/ChunkMessageChannelItemWriter.java +++ b/spring-batch-integration/src/main/java/org/springframework/batch/integration/chunk/ChunkMessageChannelItemWriter.java @@ -9,7 +9,6 @@ import org.springframework.batch.core.ExitStatus; import org.springframework.batch.core.StepContribution; import org.springframework.batch.core.StepExecution; import org.springframework.batch.core.listener.StepExecutionListenerSupport; -import org.springframework.batch.core.step.item.Chunk; import org.springframework.batch.item.ExecutionContext; import org.springframework.batch.item.ItemStream; import org.springframework.batch.item.ItemStreamException; @@ -21,12 +20,12 @@ public class ChunkMessageChannelItemWriter extends StepExecutionListenerSuppo private static final Log logger = LogFactory.getLog(ChunkMessageChannelItemWriter.class); - static final String ACTUAL = ChunkMessageChannelItemWriter.class.getName()+".ACTUAL"; + static final String ACTUAL = ChunkMessageChannelItemWriter.class.getName() + ".ACTUAL"; - static final String EXPECTED = ChunkMessageChannelItemWriter.class.getName()+".EXPECTED"; + static final String EXPECTED = ChunkMessageChannelItemWriter.class.getName() + ".EXPECTED"; private static final long DEFAULT_THROTTLE_LIMIT = 6; - + private MessagingGateway messagingGateway; private LocalState localState = new LocalState(); @@ -56,7 +55,8 @@ public class ChunkMessageChannelItemWriter extends StepExecutionListenerSuppo if (!items.isEmpty()) { logger.debug("Dispatching chunk: " + items); - ChunkRequest request = new ChunkRequest(new Chunk(items), localState.getJobId(), localState.createStepContribution()); + ChunkRequest request = new ChunkRequest(items, localState.getJobId(), localState + .createStepContribution()); messagingGateway.send(request); localState.expected++; @@ -141,8 +141,9 @@ public class ChunkMessageChannelItemWriter extends StepExecutionListenerSuppo + jobInstanceId + "] should have been [" + localState.getJobId() + "]."); localState.actual++; // TODO: apply the skip count - if (! payload.isSuccessful()) { - throw new AsynchronousFailureException("Failure or interrupt detected in handler: "+payload.getMessage()); + if (!payload.isSuccessful()) { + throw new AsynchronousFailureException("Failure or interrupt detected in handler: " + + payload.getMessage()); } } } diff --git a/spring-batch-integration/src/main/java/org/springframework/batch/integration/chunk/ChunkRequest.java b/spring-batch-integration/src/main/java/org/springframework/batch/integration/chunk/ChunkRequest.java index a6e553a93..e7fe535c2 100644 --- a/spring-batch-integration/src/main/java/org/springframework/batch/integration/chunk/ChunkRequest.java +++ b/spring-batch-integration/src/main/java/org/springframework/batch/integration/chunk/ChunkRequest.java @@ -1,6 +1,7 @@ package org.springframework.batch.integration.chunk; import java.io.Serializable; +import java.util.Collection; import org.springframework.batch.core.StepContribution; import org.springframework.batch.core.step.item.Chunk; @@ -11,8 +12,8 @@ public class ChunkRequest implements Serializable { private final Chunk items; private final StepContribution stepContribution; - public ChunkRequest(Chunk items, Long jobId, StepContribution stepContribution) { - this.items = items; + public ChunkRequest(Collection items, Long jobId, StepContribution stepContribution) { + this.items = new Chunk(items); this.jobId = jobId; this.stepContribution = stepContribution; } diff --git a/spring-batch-integration/src/test/java/org/springframework/batch/integration/chunk/ChunkMessageItemWriterIntegrationTests.java b/spring-batch-integration/src/test/java/org/springframework/batch/integration/chunk/ChunkMessageItemWriterIntegrationTests.java index 72a010e5f..ed4080671 100644 --- a/spring-batch-integration/src/test/java/org/springframework/batch/integration/chunk/ChunkMessageItemWriterIntegrationTests.java +++ b/spring-batch-integration/src/test/java/org/springframework/batch/integration/chunk/ChunkMessageItemWriterIntegrationTests.java @@ -27,7 +27,6 @@ import org.springframework.batch.core.repository.dao.MapJobExecutionDao; import org.springframework.batch.core.repository.dao.MapJobInstanceDao; import org.springframework.batch.core.repository.dao.MapStepExecutionDao; import org.springframework.batch.core.repository.support.SimpleJobRepository; -import org.springframework.batch.core.step.item.Chunk; import org.springframework.batch.core.step.item.SimpleStepFactoryBean; import org.springframework.batch.item.ExecutionContext; import org.springframework.batch.item.support.ListItemReader; @@ -192,8 +191,7 @@ public class ChunkMessageItemWriterIntegrationTests { private GenericMessage getSimpleMessage(String string, Long jobId) { StepContribution stepContribution = new JobExecution(new JobInstance(0L, new JobParameters(), "job"), 1L) .createStepExecution("step").createStepContribution(); - ChunkRequest chunk = new ChunkRequest(new Chunk(StringUtils.commaDelimitedListToSet(string)), jobId, - stepContribution); + ChunkRequest chunk = new ChunkRequest(StringUtils.commaDelimitedListToSet(string), jobId, stepContribution); GenericMessage message = new GenericMessage(chunk); return message; } diff --git a/spring-batch-integration/src/test/java/org/springframework/batch/integration/chunk/ChunkProcessorChunkHandlerTests.java b/spring-batch-integration/src/test/java/org/springframework/batch/integration/chunk/ChunkProcessorChunkHandlerTests.java index 40b784959..9c59f7834 100644 --- a/spring-batch-integration/src/test/java/org/springframework/batch/integration/chunk/ChunkProcessorChunkHandlerTests.java +++ b/spring-batch-integration/src/test/java/org/springframework/batch/integration/chunk/ChunkProcessorChunkHandlerTests.java @@ -4,12 +4,10 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import org.junit.Test; -import org.springframework.batch.core.JobExecution; -import org.springframework.batch.core.JobInstance; -import org.springframework.batch.core.JobParameters; import org.springframework.batch.core.StepContribution; import org.springframework.batch.core.step.item.Chunk; import org.springframework.batch.core.step.item.ChunkProcessor; +import org.springframework.batch.test.MetaDataInstanceFactory; import org.springframework.util.StringUtils; public class ChunkProcessorChunkHandlerTests { @@ -25,10 +23,10 @@ public class ChunkProcessorChunkHandlerTests { count += chunk.size(); } }); - StepContribution stepContribution = new JobExecution(new JobInstance(0L, new JobParameters(), "job"), 1L).createStepExecution("step").createStepContribution(); + StepContribution stepContribution = MetaDataInstanceFactory.createStepExecution().createStepContribution(); @SuppressWarnings("unchecked") - ChunkResponse response = handler.handleChunk(new ChunkRequest(new Chunk(StringUtils - .commaDelimitedListToSet("foo,bar")), 12L, stepContribution)); + ChunkResponse response = handler.handleChunk(new ChunkRequest(StringUtils + .commaDelimitedListToSet("foo,bar"), 12L, stepContribution)); assertEquals(stepContribution, response.getStepContribution()); assertEquals(12, response.getJobId().longValue()); assertTrue(response.isSuccessful());