diff --git a/spring-batch-integration/.settings/com.springsource.sts.config.flow.prefs b/spring-batch-integration/.settings/com.springsource.sts.config.flow.prefs
index edfc2175b..c60d82d8b 100644
--- a/spring-batch-integration/.settings/com.springsource.sts.config.flow.prefs
+++ b/spring-batch-integration/.settings/com.springsource.sts.config.flow.prefs
@@ -1,15 +1,17 @@
-#Fri Sep 17 16:25:48 BST 2010
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/ChunkStepIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepIntegrationTests-context.xml=\n\n\n\n\n\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJdbcIntegrationTests-context.xml=\n\n\n\n\n\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJmsIntegrationTests-context.xml=\r\n\r\n\r\n\r\n\r\n\r\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkStepIntegrationTests-context.xml=\n\n\n\n\n\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/file/FileToMessagesJobIntegrationTests-context.xml=\r\n\r\n\r\n\r\n\r\n\r\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/VanillaIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/step/StepGatewayIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/tasklet/StepGatewayIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJdbcIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJmsIntegrationTests-context.xml=\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n
-//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/retry/RetryTransactionalPollingIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n
-eclipse.preferences.version=1
+#Wed Mar 02 08:57:31 GMT 2011
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/ChunkStepIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepIntegrationTests-context.xml=\n\n\n\n\n\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJdbcIntegrationTests-context.xml=\n\n\n\n\n\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJmsIntegrationTests-context.xml=\r\n\r\n\r\n\r\n\r\n\r\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkStepIntegrationTests-context.xml=\n\n\n\n\n\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/file/FileToMessagesJobIntegrationTests-context.xml=\r\n\r\n\r\n\r\n\r\n\r\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/JmsIntegrationTests-context.xml=\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/VanillaIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/step/StepGatewayIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/tasklet/StepGatewayIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJdbcIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJmsIntegrationTests-context.xml=\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/JmsIntegrationTests-context.xml=\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n\r\n
+//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/integration\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/retry/RetryTransactionalPollingIntegrationTests-context.xml=\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n\n
+eclipse.preferences.version=1
diff --git a/spring-batch-integration/src/main/java/org/springframework/batch/integration/partition/MessageChannelPartitionHandler.java b/spring-batch-integration/src/main/java/org/springframework/batch/integration/partition/MessageChannelPartitionHandler.java
index 9a42cc4ba..409a08fde 100644
--- a/spring-batch-integration/src/main/java/org/springframework/batch/integration/partition/MessageChannelPartitionHandler.java
+++ b/spring-batch-integration/src/main/java/org/springframework/batch/integration/partition/MessageChannelPartitionHandler.java
@@ -4,6 +4,8 @@ import java.util.Collection;
import java.util.List;
import java.util.Set;
+import org.apache.commons.logging.Log;
+import org.apache.commons.logging.LogFactory;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.StepExecution;
import org.springframework.batch.core.partition.PartitionHandler;
@@ -14,23 +16,20 @@ import org.springframework.integration.MessageChannel;
import org.springframework.integration.annotation.Aggregator;
import org.springframework.integration.annotation.MessageEndpoint;
import org.springframework.integration.annotation.Payloads;
+import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.core.MessagingOperations;
import org.springframework.integration.core.PollableChannel;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.util.Assert;
/**
- * A {@link PartitionHandler} that uses {@link MessageChannel} instances to send
- * instructions to remote workers and receive their responses. The
- * {@link MessageChannel} provides a nice abstraction so that the location of
- * the workers and the transport used to communicate with them can be changed at
- * run time. The communication with the remote workers does not need to be
- * transactional or have guaranteed delivery, so a local thread pool based
- * implementation works as well as a remote web service or JMS implementation.
- * If a remote worker fails or doesn't send a reply message, the job will fail
- * and can be restarted to pick up missing messages and processing. The remote
- * workers need access to the Spring Batch {@link JobRepository} so that the
- * shared state across those restarts can be managed centrally.
+ * A {@link PartitionHandler} that uses {@link MessageChannel} instances to send instructions to remote workers and
+ * receive their responses. The {@link MessageChannel} provides a nice abstraction so that the location of the workers
+ * and the transport used to communicate with them can be changed at run time. The communication with the remote workers
+ * does not need to be transactional or have guaranteed delivery, so a local thread pool based implementation works as
+ * well as a remote web service or JMS implementation. If a remote worker fails or doesn't send a reply message, the job
+ * will fail and can be restarted to pick up missing messages and processing. The remote workers need access to the
+ * Spring Batch {@link JobRepository} so that the shared state across those restarts can be managed centrally.
*
* @author Dave Syer
*
@@ -38,31 +37,25 @@ import org.springframework.util.Assert;
@MessageEndpoint
public class MessageChannelPartitionHandler implements PartitionHandler {
+ private static Log logger = LogFactory.getLog(MessageChannelPartitionHandler.class);
+
private int gridSize = 1;
private MessagingOperations messagingGateway;
private String stepName;
- private PollableChannel replyChannel;
-
public void afterPropertiesSet() throws Exception {
Assert.notNull(stepName, "A step name must be provided for the remote workers.");
Assert.state(messagingGateway != null, "The MessagingOperations must be set");
}
/**
- * A pre-configured gateway for sending and receiving messages to the remote
- * workers. Using this property allows a large degree of control over the
- * timeouts and other properties of the send. It should have channels set up
- * internally:
- *
- * - request channel capable of accepting {@link StepExecutionRequest}
- * payloads
- * - reply channel that returns a list of {@link StepExecution} results
- *
- * The timeout for the repoy should be set sufficiently long that the remote
- * steps have time to complete.
+ * A pre-configured gateway for sending and receiving messages to the remote workers. Using this property allows a
+ * large degree of control over the timeouts and other properties of the send. It should have channels set up
+ * internally: - request channel capable of accepting {@link StepExecutionRequest} payloads
- reply
+ * channel that returns a list of {@link StepExecution} results
The timeout for the repoy should be set
+ * sufficiently long that the remote steps have time to complete.
*
* @param messagingGateway the {@link MessagingOperations} to set
*/
@@ -70,16 +63,10 @@ public class MessageChannelPartitionHandler implements PartitionHandler {
this.messagingGateway = messagingGateway;
}
- public void setReplyChannel(PollableChannel replyChannel) {
- this.replyChannel = replyChannel;
- }
-
/**
- * Passed to the {@link StepExecutionSplitter} in the
- * {@link #handle(StepExecutionSplitter, StepExecution)} method, instructing
- * it how many {@link StepExecution} instances are required, ideally. The
- * {@link StepExecutionSplitter} is allowed to ignore the grid size in the
- * case of a restart, since the input data partitions must be preserved.
+ * Passed to the {@link StepExecutionSplitter} in the {@link #handle(StepExecutionSplitter, StepExecution)} method,
+ * instructing it how many {@link StepExecution} instances are required, ideally. The {@link StepExecutionSplitter}
+ * is allowed to ignore the grid size in the case of a restart, since the input data partitions must be preserved.
*
* @param gridSize the number of step executions that will be created
*/
@@ -88,14 +75,12 @@ public class MessageChannelPartitionHandler implements PartitionHandler {
}
/**
- * The name of the {@link Step} that will be used to execute the partitioned
- * {@link StepExecution}. This is a regular Spring Batch step, with all the
- * business logic required to complete an execution based on the input
- * parameters in its {@link StepExecution} context. The name will be
- * translated into a {@link Step} instance by the remote worker.
+ * The name of the {@link Step} that will be used to execute the partitioned {@link StepExecution}. This is a
+ * regular Spring Batch step, with all the business logic required to complete an execution based on the input
+ * parameters in its {@link StepExecution} context. The name will be translated into a {@link Step} instance by the
+ * remote worker.
*
- * @param stepName the name of the {@link Step} instance to execute business
- * logic
+ * @param stepName the name of the {@link Step} instance to execute business logic
*/
public void setStepName(String stepName) {
this.stepName = stepName;
@@ -111,13 +96,10 @@ public class MessageChannelPartitionHandler implements PartitionHandler {
}
/**
- * Sends {@link StepExecutionRequest} objects to the request channel of the
- * {@link MessagingOperations}, and then receives the result back as a list of
- * {@link StepExecution} on a reply channel. Use the
- * {@link #aggregate(List)} method as an aggregator of the individual remote
- * replies. The receive timeout needs to be set realistically in the
- * {@link MessagingOperations} and the aggregator, so that there is a
- * good chance of all work being done.
+ * Sends {@link StepExecutionRequest} objects to the request channel of the {@link MessagingOperations}, and then
+ * receives the result back as a list of {@link StepExecution} on a reply channel. Use the {@link #aggregate(List)}
+ * method as an aggregator of the individual remote replies. The receive timeout needs to be set realistically in
+ * the {@link MessagingOperations} and the aggregator, so that there is a good chance of all work being done.
*
* @see PartitionHandler#handle(StepExecutionSplitter, StepExecution)
*/
@@ -126,21 +108,32 @@ public class MessageChannelPartitionHandler implements PartitionHandler {
Set split = stepExecutionSplitter.split(masterStepExecution, gridSize);
int count = 0;
+ PollableChannel replyChannel = new QueueChannel();
+
for (StepExecution stepExecution : split) {
- messagingGateway.send(createMessage(count++, split.size(), new StepExecutionRequest(stepName, stepExecution
- .getJobExecutionId(), stepExecution.getId())));
+ Message request = createMessage(count++, split.size(), new StepExecutionRequest(
+ stepName, stepExecution.getJobExecutionId(), stepExecution.getId()), replyChannel);
+ if (logger.isDebugEnabled()) {
+ logger.debug("Sending request: " + request);
+ }
+ messagingGateway.send(request);
}
Message> message = messagingGateway.receive(replyChannel);
+ if (logger.isDebugEnabled()) {
+ logger.debug("Received replies: " + message);
+ }
Collection result = message.getPayload();
return result;
}
private Message createMessage(int sequenceNumber, int sequenceSize,
- StepExecutionRequest stepExecutionRequest) {
- return MessageBuilder.withPayload(stepExecutionRequest).setSequenceNumber(sequenceNumber).setSequenceSize(
- sequenceSize).setCorrelationId(
- stepExecutionRequest.getJobExecutionId() + ":" + stepExecutionRequest.getStepName()).build();
+ StepExecutionRequest stepExecutionRequest, PollableChannel replyChannel) {
+ return MessageBuilder.withPayload(stepExecutionRequest).setSequenceNumber(sequenceNumber)
+ .setSequenceSize(sequenceSize)
+ .setCorrelationId(stepExecutionRequest.getJobExecutionId() + ":" + stepExecutionRequest.getStepName())
+ .setReplyChannel(replyChannel)
+ .build();
}
}
diff --git a/spring-batch-integration/src/test/java/org/springframework/batch/integration/partition/JmsIntegrationTests.java b/spring-batch-integration/src/test/java/org/springframework/batch/integration/partition/JmsIntegrationTests.java
new file mode 100755
index 000000000..d22532868
--- /dev/null
+++ b/spring-batch-integration/src/test/java/org/springframework/batch/integration/partition/JmsIntegrationTests.java
@@ -0,0 +1,79 @@
+/*
+ * Copyright 2006-2007 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.batch.integration.partition;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotNull;
+
+import java.util.List;
+
+import org.apache.commons.logging.Log;
+import org.apache.commons.logging.LogFactory;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.springframework.batch.core.BatchStatus;
+import org.springframework.batch.core.Job;
+import org.springframework.batch.core.JobExecution;
+import org.springframework.batch.core.JobInstance;
+import org.springframework.batch.core.JobParameters;
+import org.springframework.batch.core.StepExecution;
+import org.springframework.batch.core.explore.JobExplorer;
+import org.springframework.batch.core.launch.JobLauncher;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.test.context.ContextConfiguration;
+import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
+
+/**
+ * @author Dave Syer
+ *
+ */
+@ContextConfiguration
+@RunWith(SpringJUnit4ClassRunner.class)
+public class JmsIntegrationTests {
+
+ private Log logger = LogFactory.getLog(getClass());
+
+ @Autowired
+ private JobLauncher jobLauncher;
+
+ @Autowired
+ private Job job;
+
+ @Autowired
+ private JobExplorer jobExplorer;
+
+ @Test
+ public void testSimpleProperties() throws Exception {
+ assertNotNull(jobLauncher);
+ }
+
+ @Test
+ public void testLaunchJob() throws Exception {
+ int before = jobExplorer.getJobInstances(job.getName(), 0, 100).size();
+ assertNotNull(jobLauncher.run(job, new JobParameters()));
+ List jobInstances = jobExplorer.getJobInstances(job.getName(), 0, 100);
+ int after = jobInstances.size();
+ assertEquals(1, after - before);
+ JobExecution jobExecution = jobExplorer.getJobExecutions(jobInstances.get(jobInstances.size() - 1)).get(0);
+ assertEquals(BatchStatus.COMPLETED, jobExecution.getStatus());
+ assertEquals(3, jobExecution.getStepExecutions().size());
+ for (StepExecution stepExecution : jobExecution.getStepExecutions()) {
+ // BATCH-1703: we are using a map dao so the step executions in the job execution are old and we need to
+ // pull them back out of the repository...
+ stepExecution = jobExplorer.getStepExecution(jobExecution.getId(), stepExecution.getId());
+ logger.debug("" + stepExecution);
+ assertEquals(BatchStatus.COMPLETED, stepExecution.getStatus());
+ }
+ }
+
+}
diff --git a/spring-batch-integration/src/test/resources/jms-context.xml b/spring-batch-integration/src/test/resources/jms-context.xml
index 2d86ef606..8bbca6eb1 100644
--- a/spring-batch-integration/src/test/resources/jms-context.xml
+++ b/spring-batch-integration/src/test/resources/jms-context.xml
@@ -23,7 +23,7 @@
-
+
vm://localhost
diff --git a/spring-batch-integration/src/test/resources/log4j.properties b/spring-batch-integration/src/test/resources/log4j.properties
index 20474e2a4..fc0e914ff 100644
--- a/spring-batch-integration/src/test/resources/log4j.properties
+++ b/spring-batch-integration/src/test/resources/log4j.properties
@@ -9,4 +9,5 @@ log4j.category.org.springframework.beans=INFO
log4j.category.org.springframework.batch.retry=DEBUG
log4j.category.org.springframework.batch.core=DEBUG
log4j.category.org.springframework.batch.integration=DEBUG
+log4j.category.org.springframework.integration=DEBUG
log4j.category.org.springframework.transaction=INFO
\ No newline at end of file
diff --git a/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/JmsIntegrationTests-context.xml b/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/JmsIntegrationTests-context.xml
new file mode 100755
index 000000000..69765e57b
--- /dev/null
+++ b/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/JmsIntegrationTests-context.xml
@@ -0,0 +1,80 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/VanillaIntegrationTests-context.xml b/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/VanillaIntegrationTests-context.xml
index 2cb14174c..20b49db58 100644
--- a/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/VanillaIntegrationTests-context.xml
+++ b/spring-batch-integration/src/test/resources/org/springframework/batch/integration/partition/VanillaIntegrationTests-context.xml
@@ -15,9 +15,6 @@
-
-
-
@@ -25,7 +22,7 @@
-
@@ -49,7 +46,6 @@
-