RESOLVED - issue BATCH-814: JobRepository should not require Step or Job (only their names)

Moved the check for restartability up into JobLauncher.  There is a corner case where a non-restartable job could be restarted (see comment in SimpleJobLancher), but it should have a pretty small impact.
This commit is contained in:
dsyer
2008-09-05 17:21:24 +00:00
parent db608b91a8
commit 34c2fce32f
26 changed files with 152 additions and 128 deletions

View File

@@ -81,7 +81,17 @@ public class SimpleJobLauncher implements JobLauncher, InitializingBean {
Assert.notNull(job, "The Job must not be null.");
Assert.notNull(jobParameters, "The JobParameters must not be null.");
final JobExecution jobExecution = jobRepository.createJobExecution(job, jobParameters);
boolean exists = jobRepository.isJobInstanceExists(job.getName(), jobParameters);
if (exists && !job.isRestartable()) {
throw new JobRestartException("JobInstance already exists and is not restartable");
}
/**
* There is a very small probability that a non-restartable job can be
* restarted, but only if another process or thread manages to launch
* <i>and</i> fail a job execution for this instance between the last assertion
* and the next method returning successfully.
*/
final JobExecution jobExecution = jobRepository.createJobExecution(job.getName(), jobParameters);
taskExecutor.execute(new Runnable() {

View File

@@ -85,7 +85,8 @@ public class SimpleJobOperator implements JobOperator, InitializingBean {
public void afterPropertiesSet() throws Exception {
Assert.notNull(jobLauncher, "JobLauncher must be provided");
Assert.notNull(jobRegistry, "JobLocator must be provided");
Assert.notNull(jobExplorer, "BatchMetaDataExplorer must be provided");
Assert.notNull(jobExplorer, "JobExplorer must be provided");
Assert.notNull(jobRepository, "JobRepository must be provided");
}
/**

View File

@@ -54,9 +54,8 @@ public interface JobRepository {
* returned in a new {@link JobInstance}, associated with the
* {@link JobExecution}. If no previous instance is found, the execution
* will be associated with a new {@link JobInstance}
*
* @param jobName TODO
* @param jobParameters the runtime parameters for the job
* @param job the job the execution should be associated with.
*
* @return a valid job {@link JobExecution} for the arguments provided
* @throws JobExecutionAlreadyRunningException if there is a
@@ -69,7 +68,7 @@ public interface JobRepository {
* found and was already completed successfully.
*
*/
JobExecution createJobExecution(Job job, JobParameters jobParameters) throws JobExecutionAlreadyRunningException,
JobExecution createJobExecution(String jobName, JobParameters jobParameters) throws JobExecutionAlreadyRunningException,
JobRestartException, JobInstanceAlreadyCompleteException;
/**

View File

@@ -83,7 +83,7 @@ public class JobRepositoryFactoryBean extends AbstractJobRepositoryFactoryBean i
*
* @param isolationLevelForCreate the isolation level name to set
*
* @see SimpleJobRepository#createJobExecution(org.springframework.batch.core.Job,
* @see SimpleJobRepository#createJobExecution(String,
* org.springframework.batch.core.JobParameters)
*/
public void setIsolationLevelForCreate(String isolationLevelForCreate) {

View File

@@ -143,13 +143,13 @@ public class SimpleJobRepository implements JobRepository {
* support the higher isolation levels).
* </p>
*
* @see JobRepository#createJobExecution(Job, JobParameters)
* @see JobRepository#createJobExecution(String, JobParameters)
*
*/
public JobExecution createJobExecution(Job job, JobParameters jobParameters)
public JobExecution createJobExecution(String jobName, JobParameters jobParameters)
throws JobExecutionAlreadyRunningException, JobRestartException, JobInstanceAlreadyCompleteException {
Assert.notNull(job, "Job must not be null.");
Assert.notNull(jobName, "Job name must not be null.");
Assert.notNull(jobParameters, "JobParameters must not be null.");
/*
@@ -161,14 +161,11 @@ public class SimpleJobRepository implements JobRepository {
* has finished.
*/
JobInstance jobInstance = jobInstanceDao.getJobInstance(job.getName(), jobParameters);
JobInstance jobInstance = jobInstanceDao.getJobInstance(jobName, jobParameters);
ExecutionContext executionContext;
// existing job instance found
if (jobInstance != null) {
if (!job.isRestartable()) {
throw new JobRestartException("JobInstance already exists and is not restartable");
}
List<JobExecution> executions = jobExecutionDao.findJobExecutions(jobInstance);
@@ -188,7 +185,7 @@ public class SimpleJobRepository implements JobRepository {
}
else {
// no job found, create one
jobInstance = jobInstanceDao.createJobInstance(job.getName(), jobParameters);
jobInstance = jobInstanceDao.createJobInstance(jobName, jobParameters);
executionContext = new ExecutionContext();
}

View File

@@ -142,7 +142,7 @@ public class JobSupport implements BeanNameAware, Job {
*/
public void execute(JobExecution execution) throws UnexpectedJobExecutionException {
throw new UnsupportedOperationException(
"JobSupport does not provide an implementation of run(). Use a smarter subclass.");
"JobSupport does not provide an implementation of execute(). Use a smarter subclass.");
}
public String toString() {

View File

@@ -123,7 +123,7 @@ public class SimpleJobTests extends TestCase {
job.setName("testJob");
job.setSteps(steps);
jobExecution = jobRepository.createJobExecution(job, jobParameters);
jobExecution = jobRepository.createJobExecution(job.getName(), jobParameters);
jobInstance = jobExecution.getJobInstance();
stepExecution1 = new StepExecution(step1.getName(), jobExecution);
@@ -403,7 +403,7 @@ public class SimpleJobTests extends TestCase {
assertFalse(jobExecution.getExecutionContext().isEmpty());
jobExecution = jobRepository.createJobExecution(job, jobParameters);
jobExecution = jobRepository.createJobExecution(job.getName(), jobParameters);
try {
job.execute(jobExecution);

View File

@@ -16,8 +16,14 @@
package org.springframework.batch.core.launch;
import static org.easymock.EasyMock.*;
import static org.junit.Assert.*;
import static org.easymock.EasyMock.createMock;
import static org.easymock.EasyMock.expect;
import static org.easymock.EasyMock.replay;
import static org.easymock.EasyMock.reset;
import static org.easymock.EasyMock.verify;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
import java.util.ArrayList;
import java.util.List;
@@ -30,6 +36,7 @@ import org.springframework.batch.core.JobParameters;
import org.springframework.batch.core.job.JobSupport;
import org.springframework.batch.core.launch.support.SimpleJobLauncher;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.repository.JobRestartException;
import org.springframework.batch.repeat.ExitStatus;
import org.springframework.core.task.TaskExecutor;
@@ -42,6 +49,7 @@ public class SimpleJobLauncherTests {
private SimpleJobLauncher jobLauncher;
private Job job = new JobSupport("foo") {
@Override
public void execute(JobExecution execution) {
execution.setExitStatus(ExitStatus.FINISHED);
return;
@@ -66,8 +74,8 @@ public class SimpleJobLauncherTests {
JobExecution jobExecution = new JobExecution(null, null);
expect(jobRepository.createJobExecution(job, jobParameters)).andReturn(jobExecution);
expect(jobRepository.isJobInstanceExists(job.getName(), jobParameters)).andReturn(false);
expect(jobRepository.createJobExecution(job.getName(), jobParameters)).andReturn(jobExecution);
replay(jobRepository);
jobLauncher.afterPropertiesSet();
@@ -77,6 +85,38 @@ public class SimpleJobLauncherTests {
verify(jobRepository);
}
/*
* Non-restartable JobInstance can be run only once - attempt to run
* existing non-restartable JobInstance causes error.
*/
@Test
public void testRunNonRestartableJobInstanceTwice() throws Exception {
job = new JobSupport("foo") {
@Override
public boolean isRestartable() {
return false;
}
@Override
public void execute(JobExecution execution) {
execution.setExitStatus(ExitStatus.FINISHED);
return;
}
};
testRun();
try {
reset(jobRepository);
expect(jobRepository.isJobInstanceExists(job.getName(), jobParameters)).andReturn(true);
replay(jobRepository);
jobLauncher.run(job, jobParameters);
fail("Expected JobRestartException");
}
catch (JobRestartException e) {
// expected
}
verify(jobRepository);
}
@Test
public void testTaskExecutor() throws Exception {
final List<String> list = new ArrayList<String>();

View File

@@ -32,7 +32,6 @@ import org.springframework.batch.core.JobInstance;
import org.springframework.batch.core.JobParameters;
import org.springframework.batch.core.Step;
import org.springframework.batch.core.StepExecution;
import org.springframework.batch.core.job.JobSupport;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.step.StepSupport;
import org.springframework.dao.OptimisticLockingFailureException;
@@ -71,9 +70,7 @@ public abstract class AbstractStepExecutionDaoTests extends AbstractTransactiona
@Before
public void onSetUp() throws Exception {
repository = getJobRepository();
jobExecution = repository.createJobExecution(new JobSupport("testJob"), new JobParameters());
jobExecution = repository.createJobExecution("job", new JobParameters());
jobInstance = jobExecution.getJobInstance();
step = new StepSupport("foo");
stepExecution = new StepExecution(step.getName(), jobExecution);

View File

@@ -18,7 +18,6 @@ package org.springframework.batch.core.repository.support;
import static junit.framework.Assert.*;
import static org.easymock.EasyMock.*;
import org.springframework.batch.core.JobParameters;
import org.springframework.batch.core.job.JobSupport;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.item.database.support.DataFieldMaxValueIncrementerFactory;
import org.springframework.dao.DataAccessException;
@@ -161,7 +160,7 @@ public class JobRepositoryFactoryBeanTests {
expect(transactionManager.getTransaction(transactionDefinition)).andReturn(null);
replay(transactionManager);
try {
repository.createJobExecution(new JobSupport("job"), new JobParameters());
repository.createJobExecution("foo", new JobParameters());
// we expect an exception from the txControl because we provided the
// wrong meta data
fail("Expected IllegalArgumentException");
@@ -188,7 +187,7 @@ public class JobRepositoryFactoryBeanTests {
replay(dataSource);
replay(transactionManager);
try {
repository.createJobExecution(new JobSupport("job"), new JobParameters());
repository.createJobExecution("foo", new JobParameters());
// we expect an exception but not from the txControl because we
// provided the correct meta data
fail("Expected IllegalArgumentException");
@@ -215,7 +214,7 @@ public class JobRepositoryFactoryBeanTests {
replay(dataSource);
replay(transactionManager);
try {
repository.createJobExecution(new JobSupport("job"), new JobParameters());
repository.createJobExecution("foo", new JobParameters());
// we expect an exception but not from the txControl because we
// provided the correct meta data
fail("Expected IllegalArgumentException");

View File

@@ -1,11 +1,13 @@
package org.springframework.batch.core.repository.support;
import static org.junit.Assert.fail;
import org.junit.Test;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.JobParameters;
import org.springframework.batch.core.job.JobSupport;
import org.springframework.batch.core.repository.JobExecutionAlreadyRunningException;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.core.repository.JobRestartException;
/**
* Tests for {@link MapJobRepositoryFactoryBean}.
@@ -24,12 +26,13 @@ public class MapJobRepositoryFactoryBeanTests {
Job job = new JobSupport("jobName");
JobParameters jobParameters = new JobParameters();
repository.createJobExecution(job, jobParameters);
repository.createJobExecution(job.getName(), jobParameters);
try {
repository.createJobExecution(job, jobParameters);
repository.createJobExecution(job.getName(), jobParameters);
fail("Expected JobExecutionAlreadyRunningException");
}
catch (JobRestartException e) {
catch (JobExecutionAlreadyRunningException e) {
// expected
}
}

View File

@@ -1,9 +1,13 @@
package org.springframework.batch.core.repository.support;
import static org.junit.Assert.*;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.fail;
import java.util.Date;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.JobParameters;
@@ -12,15 +16,12 @@ import org.springframework.batch.core.Step;
import org.springframework.batch.core.StepExecution;
import org.springframework.batch.core.job.JobSupport;
import org.springframework.batch.core.repository.JobExecutionAlreadyRunningException;
import org.springframework.batch.core.repository.JobRestartException;
import org.springframework.batch.core.step.StepSupport;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import org.springframework.transaction.annotation.Transactional;
import org.junit.runner.RunWith;
import org.junit.Test;
/**
* Repository tests using JDBC DAOs (rather than mocks).
@@ -55,7 +56,7 @@ public class SimpleJobRepositoryIntegrationTests {
builder.addString("stringKey", "stringValue").addLong("longKey", 1L).addDouble("doubleKey", 1.1).addDate("dateKey", new Date(1L));
JobParameters jobParams = builder.toJobParameters();
JobExecution firstExecution = jobRepository.createJobExecution(job, jobParams);
JobExecution firstExecution = jobRepository.createJobExecution(job.getName(), jobParams);
firstExecution.setStartTime(new Date());
assertNotNull(firstExecution.getLastUpdated());
@@ -64,7 +65,7 @@ public class SimpleJobRepositoryIntegrationTests {
jobRepository.update(firstExecution);
firstExecution.setEndTime(new Date());
jobRepository.update(firstExecution);
JobExecution secondExecution = jobRepository.createJobExecution(job, jobParams);
JobExecution secondExecution = jobRepository.createJobExecution(job.getName(), jobParams);
assertEquals(firstExecution.getJobInstance(), secondExecution.getJobInstance());
assertEquals(job.getName(), secondExecution.getJobInstance().getJobName());
@@ -78,36 +79,16 @@ public class SimpleJobRepositoryIntegrationTests {
public void testCreateAndFindWithNoStartDate() throws Exception {
job.setRestartable(true);
JobExecution firstExecution = jobRepository.createJobExecution(job, jobParameters);
JobExecution firstExecution = jobRepository.createJobExecution(job.getName(), jobParameters);
firstExecution.setStartTime(new Date(0));
firstExecution.setEndTime(new Date(1));
jobRepository.update(firstExecution);
JobExecution secondExecution = jobRepository.createJobExecution(job, jobParameters);
JobExecution secondExecution = jobRepository.createJobExecution(job.getName(), jobParameters);
assertEquals(firstExecution.getJobInstance(), secondExecution.getJobInstance());
assertEquals(job.getName(), secondExecution.getJobInstance().getJobName());
}
/*
* Non-restartable JobInstance can be run only once - attempt to run
* existing non-restartable JobInstance causes error.
*/
@Transactional @Test
public void testRunNonRestartableJobInstanceTwice() throws Exception {
job.setRestartable(false);
JobExecution firstExecution = jobRepository.createJobExecution(job, jobParameters);
jobRepository.update(firstExecution);
try {
jobRepository.createJobExecution(job, jobParameters);
fail();
}
catch (JobRestartException e) {
// expected
}
}
/*
* Save multiple StepExecutions for the same step and check the returned
* count and last execution are correct.
@@ -118,7 +99,7 @@ public class SimpleJobRepositoryIntegrationTests {
StepSupport step = new StepSupport("restartedStep");
// first execution
JobExecution firstJobExec = jobRepository.createJobExecution(job, jobParameters);
JobExecution firstJobExec = jobRepository.createJobExecution(job.getName(), jobParameters);
StepExecution firstStepExec = new StepExecution(step.getName(), firstJobExec);
jobRepository.update(firstJobExec);
jobRepository.add(firstStepExec);
@@ -137,7 +118,7 @@ public class SimpleJobRepositoryIntegrationTests {
jobRepository.update(firstJobExec);
// second execution
JobExecution secondJobExec = jobRepository.createJobExecution(job, jobParameters);
JobExecution secondJobExec = jobRepository.createJobExecution(job.getName(), jobParameters);
StepExecution secondStepExec = new StepExecution(step.getName(), secondJobExec);
jobRepository.update(secondJobExec);
jobRepository.add(secondStepExec);
@@ -156,7 +137,7 @@ public class SimpleJobRepositoryIntegrationTests {
putLong("crashedPosition", 7);
}
};
JobExecution jobExec = jobRepository.createJobExecution(job, jobParameters);
JobExecution jobExec = jobRepository.createJobExecution(job.getName(), jobParameters);
jobExec.setStartTime(new Date(0));
jobExec.setExecutionContext(ctx);
Step step = new StepSupport("step1");
@@ -182,10 +163,10 @@ public class SimpleJobRepositoryIntegrationTests {
@Transactional @Test
public void testOnlyOneJobExecutionAllowedRunning() throws Exception {
job.setRestartable(true);
jobRepository.createJobExecution(job, jobParameters);
jobRepository.createJobExecution(job.getName(), jobParameters);
try {
jobRepository.createJobExecution(job, jobParameters);
jobRepository.createJobExecution(job.getName(), jobParameters);
fail();
}
catch (JobExecutionAlreadyRunningException e) {

View File

@@ -15,7 +15,6 @@
*/
package org.springframework.batch.core.step;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.JobInstance;
import org.springframework.batch.core.JobParameters;
@@ -31,7 +30,7 @@ public class JobRepositorySupport implements JobRepository {
/* (non-Javadoc)
* @see org.springframework.batch.container.common.repository.JobRepository#findOrCreateJob(org.springframework.batch.container.common.domain.JobConfiguration)
*/
public JobExecution createJobExecution(Job jobConfiguration, JobParameters jobParameters) {
public JobExecution createJobExecution(String jobName, JobParameters jobParameters) {
return null;
}

View File

@@ -96,7 +96,7 @@ public class SimpleStepFactoryBeanTests {
step.setName("step2");
job.addStep(step);
JobExecution jobExecution = repository.createJobExecution(job, new JobParameters());
JobExecution jobExecution = repository.createJobExecution(job.getName(), new JobParameters());
job.execute(jobExecution);
assertEquals(BatchStatus.COMPLETED, jobExecution.getStatus());
@@ -116,7 +116,7 @@ public class SimpleStepFactoryBeanTests {
step.setName("step1");
job.addStep(step);
JobExecution jobExecution = repository.createJobExecution(job, new JobParameters());
JobExecution jobExecution = repository.createJobExecution(job.getName(), new JobParameters());
job.execute(jobExecution);
assertEquals(BatchStatus.COMPLETED, jobExecution.getStatus());
@@ -150,7 +150,7 @@ public class SimpleStepFactoryBeanTests {
job.setSteps(Collections.singletonList(step));
JobExecution jobExecution = repository.createJobExecution(job, new JobParameters());
JobExecution jobExecution = repository.createJobExecution(job.getName(), new JobParameters());
try {
job.execute(jobExecution);
fail("Expected RuntimeException");
@@ -179,7 +179,7 @@ public class SimpleStepFactoryBeanTests {
AbstractStep step = (AbstractStep) factory.getObject();
job.setSteps(Collections.singletonList((Step) step));
JobExecution jobExecution = repository.createJobExecution(job, new JobParameters());
JobExecution jobExecution = repository.createJobExecution(job.getName(), new JobParameters());
try {
job.execute(jobExecution);
fail("Expected RuntimeException");
@@ -210,7 +210,7 @@ public class SimpleStepFactoryBeanTests {
AbstractStep step = (AbstractStep) factory.getObject();
job.setSteps(Collections.singletonList((Step) step));
JobExecution jobExecution = repository.createJobExecution(job, new JobParameters());
JobExecution jobExecution = repository.createJobExecution(job.getName(), new JobParameters());
job.execute(jobExecution);
@@ -244,7 +244,7 @@ public class SimpleStepFactoryBeanTests {
job.setSteps(Collections.singletonList((Step) step));
JobExecution jobExecution = repository.createJobExecution(job, new JobParameters());
JobExecution jobExecution = repository.createJobExecution(job.getName(), new JobParameters());
job.execute(jobExecution);
assertEquals(BatchStatus.COMPLETED, jobExecution.getStatus());

View File

@@ -116,7 +116,7 @@ public class StatefulRetryStepFactoryBeanTests {
job.setRestartable(true);
JobParameters jobParameters = new JobParametersBuilder().addString("statefulTest", "make_this_unique")
.toJobParameters();
jobExecution = repository.createJobExecution(job, jobParameters);
jobExecution = repository.createJobExecution(job.getName(), jobParameters);
jobExecution.setEndTime(new Date());
}

View File

@@ -124,7 +124,7 @@ public class ChunkOrientedStepIntegrationTests {
}
}, chunkOperations));
JobExecution jobExecution = jobRepository.createJobExecution(job, new JobParameters());
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(), new JobParameters());
StepExecution stepExecution = new StepExecution(step.getName(), jobExecution);
stepExecution.setExecutionContext(new ExecutionContext() {

View File

@@ -58,11 +58,11 @@ public class StepExecutorInterruptionTests extends TestCase {
JobRepository jobRepository = new SimpleJobRepository(new MapJobInstanceDao(), new MapJobExecutionDao(),
new MapStepExecutionDao(), new MapExecutionContextDao());
JobSupport jobConfiguration = new JobSupport();
JobSupport job = new JobSupport();
step = new TaskletStep("interruptedStep");
jobConfiguration.addStep(step);
jobConfiguration.setBeanName("testJob");
jobExecution = jobRepository.createJobExecution(jobConfiguration, new JobParameters());
job.addStep(step);
job.setBeanName("testJob");
jobExecution = jobRepository.createJobExecution(job.getName(), new JobParameters());
step.setJobRepository(jobRepository);
step.setTransactionManager(new ResourcelessTransactionManager());
itemWriter = new ItemWriter<Object>() {

View File

@@ -192,7 +192,7 @@ public class TasketStepTests {
new MapStepExecutionDao(), new MapExecutionContextDao());
step.setJobRepository(repository);
JobExecution jobExecution = repository.createJobExecution(job, jobInstance.getJobParameters());
JobExecution jobExecution = repository.createJobExecution(job.getName(), jobInstance.getJobParameters());
StepExecution stepExecution = new StepExecution(step.getName(), jobExecution);
step.execute(stepExecution);

View File

@@ -15,7 +15,6 @@
*/
package org.springframework.batch.integration;
import org.springframework.batch.core.Job;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.JobInstance;
import org.springframework.batch.core.JobParameters;
@@ -34,7 +33,7 @@ public class JobRepositorySupport implements JobRepository {
/* (non-Javadoc)
* @see org.springframework.batch.core.repository.JobRepository#createJobExecution(org.springframework.batch.core.Job, org.springframework.batch.core.JobParameters)
*/
public JobExecution createJobExecution(Job job, JobParameters jobParameters)
public JobExecution createJobExecution(String jobName, JobParameters jobParameters)
throws JobExecutionAlreadyRunningException, JobRestartException, JobInstanceAlreadyCompleteException {
return new JobExecution(new JobInstance(0L, jobParameters, job.getName()));
}

View File

@@ -321,7 +321,7 @@ public class ChunkMessageItemWriterIntegrationTests {
JobInstanceAlreadyCompleteException {
SimpleJob job = new SimpleJob();
job.setName("job");
JobExecution jobExecution = jobRepository.createJobExecution(job,
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(),
new JobParametersBuilder().addLong("job.counter", jobCounter++)
.toJobParameters());
StepExecution stepExecution = jobExecution.createStepExecution(step.getName());

View File

@@ -176,7 +176,7 @@ public class FileToMessagesJobFactoryBeanTests {
Job job = (Job) factory.getObject();
JobParameters jobParameters = new JobParametersBuilder().addString(FILE_INPUT_PATH, "classpath:/log4j.properties").toJobParameters();
JobExecution jobExecution = jobRepository.createJobExecution(job, jobParameters);
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(), jobParameters);
job.execute(jobExecution);
assertNotNull(jobExecution);

View File

@@ -33,7 +33,6 @@ import org.springframework.batch.core.Step;
import org.springframework.batch.core.StepExecution;
import org.springframework.batch.core.repository.JobRepository;
import org.springframework.batch.integration.JobRepositorySupport;
import org.springframework.batch.integration.JobSupport;
import org.springframework.batch.integration.StepSupport;
import org.springframework.batch.item.ExecutionContext;
import org.springframework.beans.factory.annotation.Required;
@@ -81,7 +80,7 @@ public class StepExecutionMessageHandlerTests {
JobRepositorySupport jobRepository = new JobRepositorySupport();
StepExecutionMessageHandler handler = createHandler(jobRepository);
JobExecutionRequest message = handler.handle(new JobExecutionRequest(jobRepository.createJobExecution(
new JobSupport("job"), new JobParameters())));
job.getName(), new JobParameters())));
assertEquals(1, message.getJobExecution().getStepExecutions().size());
assertEquals(BatchStatus.COMPLETED, message.getStatus());
}
@@ -91,7 +90,7 @@ public class StepExecutionMessageHandlerTests {
JobRepositorySupport jobRepository = new JobRepositorySupport();
StepExecutionMessageHandler handler = createHandler(jobRepository);
JobExecutionRequest jobExecutionRequest = new JobExecutionRequest(jobRepository.createJobExecution(
new JobSupport("job"), new JobParameters()));
job.getName(), new JobParameters()));
jobExecutionRequest.getJobExecution().getExecutionContext().putString("foo", "bar");
JobExecutionRequest message = handler.handle(jobExecutionRequest);
assertEquals(1, message.getJobExecution().getStepExecutions().size());
@@ -104,7 +103,7 @@ public class StepExecutionMessageHandlerTests {
JobRepositorySupport jobRepository = new JobRepositorySupport();
StepExecutionMessageHandler handler = createHandler(jobRepository);
JobExecutionRequest jobExecutionRequest = new JobExecutionRequest(jobRepository.createJobExecution(
new JobSupport("job"), new JobParameters()));
job.getName(), new JobParameters()));
jobExecutionRequest.getJobExecution().getExecutionContext().putString("foo", "bar");
// The step has to add the output attribute to the context
handler.setStep(new StepSupport("step") {
@@ -123,7 +122,7 @@ public class StepExecutionMessageHandlerTests {
public void testHandleFailedJob() throws Exception {
JobRepositorySupport jobRepository = new JobRepositorySupport();
StepExecutionMessageHandler handler = createHandler(jobRepository);
JobExecution jobExecution = jobRepository.createJobExecution(new JobSupport("job"), new JobParameters());
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(), new JobParameters());
jobExecution.setStatus(BatchStatus.FAILED);
JobExecutionRequest message = handler.handle(new JobExecutionRequest(jobExecution));
assertEquals(0, message.getJobExecution().getStepExecutions().size());
@@ -158,7 +157,7 @@ public class StepExecutionMessageHandlerTests {
}
};
StepExecutionMessageHandler handler = createHandler(jobRepository);
JobExecution jobExecution = jobRepository.createJobExecution(new JobSupport("job"), new JobParameters());
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(), new JobParameters());
JobExecutionRequest message = handler.handle(new JobExecutionRequest(jobExecution));
assertNotNull(message);
assertEquals(1, jobExecution.getStepExecutions().size());
@@ -182,7 +181,7 @@ public class StepExecutionMessageHandlerTests {
}
};
StepExecutionMessageHandler handler = createHandler(jobRepository);
JobExecution jobExecution = jobRepository.createJobExecution(new JobSupport("job"), new JobParameters());
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(), new JobParameters());
JobExecutionRequest message = handler.handle(new JobExecutionRequest(jobExecution));
assertNotNull(message);
assertEquals(1, jobExecution.getStepExecutions().size());
@@ -208,7 +207,7 @@ public class StepExecutionMessageHandlerTests {
}
};
StepExecutionMessageHandler handler = createHandler(jobRepository);
JobExecution jobExecution = jobRepository.createJobExecution(new JobSupport("job"), new JobParameters());
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(), new JobParameters());
JobExecutionRequest message = handler.handle(new JobExecutionRequest(jobExecution));
assertNotNull(message);
assertEquals(1, jobExecution.getStepExecutions().size());

View File

@@ -5,16 +5,16 @@
xmlns:p="http://www.springframework.org/schema/p"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:util="http://www.springframework.org/schema/util"
xsi:schemaLocation="
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-2.0.xsd
http://www.springframework.org/schema/aop http://www.springframework.org/schema/aop/spring-aop-2.0.xsd
http://www.springframework.org/schema/tx http://www.springframework.org/schema/tx/spring-tx-2.0.xsd
http://www.springframework.org/schema/util http://www.springframework.org/schema/util/spring-util-2.0.xsd">
xsi:schemaLocation="
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-2.0.xsd
http://www.springframework.org/schema/aop http://www.springframework.org/schema/aop/spring-aop-2.0.xsd
http://www.springframework.org/schema/tx http://www.springframework.org/schema/tx/spring-tx-2.0.xsd
http://www.springframework.org/schema/util http://www.springframework.org/schema/util/spring-util-2.0.xsd">
<bean id="xmlStaxJob" parent="simpleJob">
<property name="steps">
<bean id="step1" parent="simpleStep">
<property name="itemReader">
<bean id="xmlStaxJob" parent="simpleJob">
<property name="steps">
<bean id="step1" parent="simpleStep">
<property name="itemReader">
<bean
class="org.springframework.batch.item.xml.StaxEventItemReader">
<property name="fragmentRootElementName"

View File

@@ -27,7 +27,6 @@ import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.JobParameters;
import org.springframework.batch.core.launch.JobOperator;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import org.springframework.transaction.annotation.Transactional;

View File

@@ -25,45 +25,46 @@ import org.junit.runner.RunWith;
import org.springframework.batch.core.BatchStatus;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.JobParameters;
import org.springframework.batch.core.JobParametersBuilder;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import org.springframework.transaction.annotation.Transactional;
/**
* Functional test for graceful shutdown. A batch container is started in a new thread,
* then it's stopped using {@link JobExecution#stop()}.
* Functional test for graceful shutdown. A batch container is started in a new
* thread, then it's stopped using {@link JobExecution#stop()}.
*
* @author Lucas Ward
*
*
*/
@RunWith(SpringJUnit4ClassRunner.class)
@ContextConfiguration()
public class GracefulShutdownFunctionalTests extends AbstractBatchLauncherTests {
@Transactional @Test
@Test
public void testLaunchJob() throws Exception {
final JobParameters jobParameters = new JobParameters();
final JobParameters jobParameters = new JobParametersBuilder().addLong("timestamp", System.currentTimeMillis())
.toJobParameters();
JobExecution jobExecution = getLauncher().run(getJob(), jobParameters);
Thread.sleep(1000);
assertEquals(BatchStatus.STARTED, jobExecution.getStatus());
assertTrue(jobExecution.isRunning());
jobExecution.stop();
int count = 0;
while(jobExecution.isRunning() && count <= 10){
logger.info("Checking for end time in JobExecution: count="+count);
while (jobExecution.isRunning() && count <= 10) {
logger.info("Checking for end time in JobExecution: count=" + count);
Thread.sleep(100);
count++;
}
assertFalse("Timed out waiting for job to end.", jobExecution.isRunning());
assertEquals(BatchStatus.STOPPED, jobExecution.getStatus());
}
}

View File

@@ -55,7 +55,7 @@ public class JdbcJobRepositoryTests {
private JobRepository repository;
private JobSupport jobConfiguration;
private JobSupport job;
private Set<Long> jobExecutionIds = new HashSet<Long>();
@@ -87,8 +87,8 @@ public class JdbcJobRepositoryTests {
@Before
public void onSetUpInTransaction() throws Exception {
jobConfiguration = new JobSupport("test-job");
jobConfiguration.setRestartable(true);
job = new JobSupport("test-job");
job.setRestartable(true);
simpleJdbcTemplate.update("DELETE FROM BATCH_EXECUTION_CONTEXT");
simpleJdbcTemplate.update("DELETE FROM BATCH_STEP_EXECUTION");
simpleJdbcTemplate.update("DELETE FROM BATCH_JOB_EXECUTION");
@@ -112,9 +112,9 @@ public class JdbcJobRepositoryTests {
@Transactional @Test
public void testFindOrCreateJob() throws Exception {
jobConfiguration.setName("foo");
job.setName("foo");
int before = simpleJdbcTemplate.queryForInt("SELECT COUNT(*) FROM BATCH_JOB_INSTANCE");
JobExecution execution = repository.createJobExecution(jobConfiguration, new JobParameters());
JobExecution execution = repository.createJobExecution(job.getName(), new JobParameters());
int after = simpleJdbcTemplate.queryForInt("SELECT COUNT(*) FROM BATCH_JOB_INSTANCE");
assertEquals(before + 1, after);
assertNotNull(execution.getId());
@@ -123,7 +123,7 @@ public class JdbcJobRepositoryTests {
@Transactional @Test
public void testFindOrCreateJobConcurrently() throws Exception {
jobConfiguration.setName("bar");
job.setName("bar");
int before = simpleJdbcTemplate.queryForInt("SELECT COUNT(*) FROM BATCH_JOB_INSTANCE");
assertEquals(0, before);
@@ -153,9 +153,9 @@ public class JdbcJobRepositoryTests {
@Transactional @Test
public void testFindOrCreateJobConcurrentlyWhenJobAlreadyExists() throws Exception {
jobConfiguration.setName("spam");
job.setName("spam");
JobExecution execution = repository.createJobExecution(jobConfiguration, new JobParameters());
JobExecution execution = repository.createJobExecution(job.getName(), new JobParameters());
cacheJobIds(execution);
execution.setEndTime(new Timestamp(System.currentTimeMillis()));
repository.update(execution);
@@ -196,7 +196,7 @@ public class JdbcJobRepositoryTests {
new TransactionTemplate(transactionManager).execute(new TransactionCallback() {
public Object doInTransaction(org.springframework.transaction.TransactionStatus status) {
try {
JobExecution execution = repository.createJobExecution(jobConfiguration, new JobParameters());
JobExecution execution = repository.createJobExecution(job.getName(), new JobParameters());
cacheJobIds(execution);
list.add(execution);
Thread.sleep(1000);
@@ -216,7 +216,7 @@ public class JdbcJobRepositoryTests {
}).start();
Thread.sleep(400);
JobExecution execution = repository.createJobExecution(jobConfiguration, new JobParameters());
JobExecution execution = repository.createJobExecution(job.getName(), new JobParameters());
cacheJobIds(execution);
int count = 0;