diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/job/AbstractJob.java b/spring-batch-core/src/main/java/org/springframework/batch/core/job/AbstractJob.java index d808e2d00..cf5d9294c 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/job/AbstractJob.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/job/AbstractJob.java @@ -327,6 +327,8 @@ public abstract class AbstractJob implements Job, BeanNameAware, InitializingBea else { currentStepExecution.setExecutionContext(new ExecutionContext()); } + + jobRepository.add(currentStepExecution); step.execute(currentStepExecution); diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/partition/support/SimpleStepExecutionSplitter.java b/spring-batch-core/src/main/java/org/springframework/batch/core/partition/support/SimpleStepExecutionSplitter.java index 626120db2..7ad04398a 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/partition/support/SimpleStepExecutionSplitter.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/partition/support/SimpleStepExecutionSplitter.java @@ -90,6 +90,7 @@ public class SimpleStepExecutionSplitter implements StepExecutionSplitter { boolean startable = getStartable(currentStepExecution, contexts.get(key)); if (startable) { + jobRepository.add(currentStepExecution); set.add(currentStepExecution); } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java b/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java index 6dafa5bf3..c0071face 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/step/AbstractStep.java @@ -183,7 +183,7 @@ public abstract class AbstractStep implements Step, InitializingBean, BeanNameAw UnexpectedJobExecutionException { stepExecution.setStartTime(new Date()); stepExecution.setStatus(BatchStatus.STARTED); - getJobRepository().add(stepExecution); + getJobRepository().update(stepExecution); ExitStatus exitStatus = ExitStatus.FAILED; Exception commitException = null; diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/job/SimpleJobTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/job/SimpleJobTests.java index 37ecab213..0953529bb 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/job/SimpleJobTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/job/SimpleJobTests.java @@ -499,7 +499,7 @@ public class SimpleJobTests { passedInStepContext = new ExecutionContext(stepExecution.getExecutionContext()); stepExecution.getExecutionContext().putString("stepKey", "stepValue"); stepExecution.getJobExecution().getExecutionContext().putString("jobKey", "jobValue"); - jobRepository.add(stepExecution); + jobRepository.update(stepExecution); jobRepository.updateExecutionContext(stepExecution); if (exception instanceof RuntimeException) { diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/partition/RestartIntegrationTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/partition/RestartIntegrationTests.java index 4d3a617c7..d0e162992 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/partition/RestartIntegrationTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/partition/RestartIntegrationTests.java @@ -87,6 +87,7 @@ public class RestartIntegrationTests { // Two attempts assertEquals(2, afterMaster-beforeMaster); // One failure and two successes + System.out.println(afterPartition); assertEquals(3, afterPartition-beforePartition); } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/partition/support/PartitionStepTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/partition/support/PartitionStepTests.java index 7362a0338..8ff0e6d78 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/partition/support/PartitionStepTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/partition/support/PartitionStepTests.java @@ -69,6 +69,7 @@ public class PartitionStepTests { }); step.afterPropertiesSet(); StepExecution stepExecution = new JobExecution(0L).createStepExecution("foo"); + jobRepository.add(stepExecution); step.execute(stepExecution); // one master and two workers assertEquals(3, stepExecution.getJobExecution().getStepExecutions().size()); @@ -91,6 +92,7 @@ public class PartitionStepTests { }); step.afterPropertiesSet(); StepExecution stepExecution = new JobExecution(0L).createStepExecution("foo"); + jobRepository.add(stepExecution); step.execute(stepExecution); // one master and two workers assertEquals(3, stepExecution.getJobExecution().getStepExecutions().size()); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/partition/support/SimpleStepExecutionSplitterTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/partition/support/SimpleStepExecutionSplitterTests.java index 56f5546a4..998fecf6d 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/partition/support/SimpleStepExecutionSplitterTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/partition/support/SimpleStepExecutionSplitterTests.java @@ -1,9 +1,10 @@ package org.springframework.batch.core.partition.support; -import static org.junit.Assert.assertEquals; +import static org.junit.Assert.*; import java.util.Collections; import java.util.Map; +import java.util.Set; import org.junit.Before; import org.junit.Test; @@ -37,7 +38,12 @@ public class SimpleStepExecutionSplitterTests { @Test public void testSimpleStepExecutionProviderJobRepositoryStep() throws Exception { SimpleStepExecutionSplitter provider = new SimpleStepExecutionSplitter(jobRepository, step); - assertEquals(2, provider.split(stepExecution, 2).size()); + Set execs = provider.split(stepExecution, 2); + assertEquals(2, execs.size()); + + for (StepExecution execution : execs) { + assertNotNull("step execution partition is saved", execution.getId()); + } } @Test diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/AbstractStepTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/AbstractStepTests.java index 84d4781a4..904bf2299 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/AbstractStepTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/AbstractStepTests.java @@ -132,6 +132,7 @@ public class AbstractStepTests { @Before public void setUp() throws Exception { tested.setJobRepository(repository); + repository.add(execution); } @Test diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java index 083b5a369..77ee384e7 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/FaultTolerantStepFactoryBeanRetryTests.java @@ -155,6 +155,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { Step step = (Step) factory.getObject(); StepExecution stepExecution = new StepExecution(step.getName(), jobExecution); + repository.add(stepExecution); step.execute(stepExecution); assertEquals(0, stepExecution.getSkipCount()); @@ -192,6 +193,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { Step step = (Step) factory.getObject(); StepExecution stepExecution = new StepExecution(step.getName(), jobExecution); + repository.add(stepExecution); step.execute(stepExecution); assertEquals(2, stepExecution.getSkipCount()); @@ -246,6 +248,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { AbstractStep step = (AbstractStep) factory.getObject(); step.setName("mytest"); StepExecution stepExecution = new StepExecution(step.getName(), jobExecution); + repository.add(stepExecution); step.execute(stepExecution); assertEquals(2, recovered.size()); @@ -310,6 +313,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { AbstractStep step = (AbstractStep) factory.getObject(); step.setName("mytest"); StepExecution stepExecution = new StepExecution(step.getName(), jobExecution); + repository.add(stepExecution); step.execute(stepExecution); assertEquals(2, recovered.size()); @@ -358,6 +362,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { Step step = (Step) factory.getObject(); StepExecution stepExecution = new StepExecution(step.getName(), jobExecution); + repository.add(stepExecution); step.execute(stepExecution); assertEquals(BatchStatus.FAILED, stepExecution.getStatus()); @@ -410,6 +415,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { Step step = (Step) factory.getObject(); StepExecution stepExecution = new StepExecution(step.getName(), jobExecution); + repository.add(stepExecution); step.execute(stepExecution); String message = stepExecution.getFailureExceptions().get(0).getMessage(); assertTrue("Wrong message: " + message, message.contains("Write error - planned but not skippable.")); @@ -452,6 +458,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { AbstractStep step = (AbstractStep) factory.getObject(); StepExecution stepExecution = new StepExecution(step.getName(), jobExecution); + repository.add(stepExecution); step.execute(stepExecution); assertEquals(BatchStatus.FAILED, stepExecution.getStatus()); @@ -500,6 +507,7 @@ public class FaultTolerantStepFactoryBeanRetryTests { AbstractStep step = (AbstractStep) factory.getObject(); StepExecution stepExecution = new StepExecution(step.getName(), jobExecution); + repository.add(stepExecution); step.execute(stepExecution); assertEquals(BatchStatus.FAILED, stepExecution.getStatus()); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/TaskletStepExceptionTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/TaskletStepExceptionTests.java index 583694acd..10fa309e2 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/TaskletStepExceptionTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/item/TaskletStepExceptionTests.java @@ -90,7 +90,7 @@ public class TaskletStepExceptionTests { taskletStep.execute(stepExecution); assertEquals(FAILED, stepExecution.getStatus()); assertTrue(stepExecution.getFailureExceptions().contains(exception)); - assertEquals(1, jobRepository.getUpdateCount()); + assertEquals(2, jobRepository.getUpdateCount()); } @Test @@ -106,7 +106,7 @@ public class TaskletStepExceptionTests { taskletStep.execute(stepExecution); assertEquals(FAILED, stepExecution.getStatus()); assertTrue(stepExecution.getFailureExceptions().contains(exception)); - assertEquals(1, jobRepository.getUpdateCount()); + assertEquals(2, jobRepository.getUpdateCount()); } @Test @@ -129,7 +129,7 @@ public class TaskletStepExceptionTests { taskletStep.execute(stepExecution); assertEquals(COMPLETED, stepExecution.getStatus()); assertFalse(stepExecution.getFailureExceptions().contains(exception)); - assertEquals(3, jobRepository.getUpdateCount()); + assertEquals(4, jobRepository.getUpdateCount()); } @Test @@ -149,7 +149,7 @@ public class TaskletStepExceptionTests { assertEquals(FAILED, stepExecution.getStatus()); assertTrue(stepExecution.getFailureExceptions().contains(taskletException)); assertFalse(stepExecution.getFailureExceptions().contains(exception)); - assertEquals(1, jobRepository.getUpdateCount()); + assertEquals(2, jobRepository.getUpdateCount()); } @Test @@ -167,7 +167,7 @@ public class TaskletStepExceptionTests { assertEquals(FAILED, stepExecution.getStatus()); assertTrue(stepExecution.getFailureExceptions().contains(taskletException)); assertTrue(stepExecution.getFailureExceptions().contains(exception)); - assertEquals(1, jobRepository.getUpdateCount()); + assertEquals(2, jobRepository.getUpdateCount()); } @Test @@ -200,12 +200,17 @@ public class TaskletStepExceptionTests { final RuntimeException exception = new RuntimeException(); taskletStep.setJobRepository(new UpdateCountingJobRepository() { + boolean firstCall = true; @Override public void update(StepExecution arg0) { + if (firstCall) { + firstCall = false; + return; + } throw exception; } }); - + taskletStep.execute(stepExecution); assertEquals(UNKNOWN, stepExecution.getStatus()); assertTrue(stepExecution.getFailureExceptions().contains(taskletException)); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/ChunkOrientedStepIntegrationTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/ChunkOrientedStepIntegrationTests.java index 10e63edc4..cce93a8d5 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/ChunkOrientedStepIntegrationTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/ChunkOrientedStepIntegrationTests.java @@ -130,8 +130,8 @@ public class ChunkOrientedStepIntegrationTests { put("foo", "bar"); } }); - // step.setLastExecution(stepExecution); - + + jobRepository.add(stepExecution); step.execute(stepExecution); assertEquals(BatchStatus.UNKNOWN, stepExecution.getStatus()); StepExecution lastStepExecution = jobRepository.getLastStepExecution(jobExecution.getJobInstance(), step.getName()); diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/StepExecutorInterruptionTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/StepExecutorInterruptionTests.java index 6f4e4a2d3..5b67f8361 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/StepExecutorInterruptionTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/StepExecutorInterruptionTests.java @@ -49,13 +49,15 @@ public class StepExecutorInterruptionTests extends TestCase { private ItemWriter itemWriter; private StepExecution stepExecution; + + private JobRepository jobRepository; public void setUp() throws Exception { MapJobInstanceDao.clear(); MapJobExecutionDao.clear(); MapStepExecutionDao.clear(); - JobRepository jobRepository = new SimpleJobRepository(new MapJobInstanceDao(), new MapJobExecutionDao(), + jobRepository = new SimpleJobRepository(new MapJobInstanceDao(), new MapJobExecutionDao(), new MapStepExecutionDao(), new MapExecutionContextDao()); JobSupport job = new JobSupport(); @@ -169,6 +171,7 @@ public class StepExecutorInterruptionTests extends TestCase { } }); + jobRepository.add(stepExecution); step.execute(stepExecution); assertEquals("Planned!", stepExecution.getFailureExceptions().get(0).getMessage()); @@ -182,6 +185,7 @@ public class StepExecutorInterruptionTests extends TestCase { Thread processingThread = new Thread() { public void run() { try { + jobRepository.add(stepExecution); step.execute(stepExecution); } catch (JobInterruptedException e) { diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/TaskletStepTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/TaskletStepTests.java index a3d599b3e..4007f75ee 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/TaskletStepTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/step/tasklet/TaskletStepTests.java @@ -186,7 +186,7 @@ public class TaskletStepTests { JobExecution jobExecution = repository.createJobExecution(job.getName(), jobInstance.getJobParameters()); StepExecution stepExecution = new StepExecution(step.getName(), jobExecution); - + repository.add(stepExecution); step.execute(stepExecution); assertEquals(1, processed.size()); } @@ -806,7 +806,7 @@ public class TaskletStepTests { public void update(StepExecution stepExecution) { updateCount++; if (updateCount <= 3) { - assertEquals(updateCount, stepExecution.getReadCount()); + assertEquals(updateCount, stepExecution.getReadCount() + 1); } } @@ -814,11 +814,11 @@ public class TaskletStepTests { private static class JobRepositoryFailedUpdateStub extends JobRepositorySupport { - private boolean firstCall = true; + private int called = 0; public void update(StepExecution stepExecution) { - if (firstCall) { - firstCall = false; + called++; + if (called == 3) { throw new DataAccessResourceFailureException("stub exception"); } }