diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDao.java b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDao.java index d315d0ddd..b0b400f5a 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDao.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDao.java @@ -43,6 +43,8 @@ public class JdbcJobExecutionDao extends AbstractJdbcBatchMetadataDao implements + "END_TIME, STATUS, CONTINUABLE, EXIT_CODE, EXIT_MESSAGE, VERSION, CREATE_TIME) values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"; private static final String CHECK_JOB_EXECUTION_EXISTS = "SELECT COUNT(*) FROM %PREFIX%JOB_EXECUTION WHERE JOB_EXECUTION_ID = ?"; + + private static final String GET_STATUS = "SELECT STATUS from %PREFIX%JOB_EXECUTION where JOB_EXECUTION_ID = ?"; private static final String UPDATE_JOB_EXECUTION = "UPDATE %PREFIX%JOB_EXECUTION set START_TIME = ?, END_TIME = ?, " + " STATUS = ?, CONTINUABLE = ?, EXIT_CODE = ?, EXIT_MESSAGE = ?, VERSION = ?, CREATE_TIME = ? where JOB_EXECUTION_ID = ?"; @@ -268,6 +270,12 @@ public class JdbcJobExecutionDao extends AbstractJdbcBatchMetadataDao implements return result; } + + public void synchronizeStatus(JobExecution jobExecution) { + + String status = getJdbcTemplate().queryForObject(getQuery(GET_STATUS), String.class, jobExecution.getId()); + jobExecution.setStatus(BatchStatus.valueOf(status)); + } /** * Re-usable mapper for {@link JobExecution} instances. @@ -296,5 +304,4 @@ public class JdbcJobExecutionDao extends AbstractJdbcBatchMetadataDao implements } } - } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JobExecutionDao.java b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JobExecutionDao.java index fc7f2dbc1..df3041430 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JobExecutionDao.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/JobExecutionDao.java @@ -60,5 +60,13 @@ public interface JobExecutionDao { * @return the {@link JobExecution} for given identifier. */ JobExecution getJobExecution(Long executionId); + + /** + * Because it may be possible that the status of a JobExecution is updated while running, + * the following method while synchronize only the status field. + * + * @param jobExecution to be updated. + */ + void synchronizeStatus(JobExecution jobExecution); } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/MapJobExecutionDao.java b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/MapJobExecutionDao.java index 2512e5e71..37d24a29d 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/MapJobExecutionDao.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/dao/MapJobExecutionDao.java @@ -87,4 +87,9 @@ public class MapJobExecutionDao implements JobExecutionDao { return executionsById.get(executionId); } + public void synchronizeStatus(JobExecution jobExecution) { + // TODO Auto-generated method stub + + } + } diff --git a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/support/SimpleJobRepository.java b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/support/SimpleJobRepository.java index 816df782b..3341c972d 100644 --- a/spring-batch-core/src/main/java/org/springframework/batch/core/repository/support/SimpleJobRepository.java +++ b/spring-batch-core/src/main/java/org/springframework/batch/core/repository/support/SimpleJobRepository.java @@ -239,9 +239,11 @@ public class SimpleJobRepository implements JobRepository { public void update(StepExecution stepExecution) { validateStepExecution(stepExecution); Assert.notNull(stepExecution.getId(), "StepExecution must already be saved (have an id assigned)"); + stepExecution.setLastUpdated(new Date(System.currentTimeMillis())); stepExecutionDao.updateStepExecution(stepExecution); + checkForInterruption(stepExecution); } private void validateStepExecution(StepExecution stepExecution) { @@ -304,5 +306,19 @@ public class SimpleJobRepository implements JobRepository { } return count; } + + /* + * Check to determine whether or not the JobExecution that is the parent of the provided + * StepExecution has been interrupted. If, after synchronizing the status with the database, + * the status has been updated to STOPPING, then the job has been interrupted. + * + * @param stepExecution + */ + private void checkForInterruption(StepExecution stepExecution){ + jobExecutionDao.synchronizeStatus(stepExecution.getJobExecution()); + if(stepExecution.getJobExecution().getStatus() == BatchStatus.STOPPING){ + stepExecution.setTerminateOnly(); + } + } } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/AbstractJobExecutionDaoTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/AbstractJobExecutionDaoTests.java index 733f97b4f..51d55afb6 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/AbstractJobExecutionDaoTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/AbstractJobExecutionDaoTests.java @@ -1,6 +1,7 @@ package org.springframework.batch.core.repository.dao; import static org.junit.Assert.*; + import org.junit.Before; import org.junit.Test; @@ -196,4 +197,15 @@ public abstract class AbstractJobExecutionDaoTests extends AbstractTransactional JobExecution value = dao.getJobExecution(54321L); assertNull(value); } + + @Transactional + @Test + public void testUpdateExecutionStatus(){ + + dao.saveJobExecution(execution); + execution.setStatus(BatchStatus.COMPLETED); + dao.synchronizeStatus(execution); + assertEquals(BatchStatus.STARTING, execution.getStatus()); + } + } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDaoTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDaoTests.java index c6e4dd348..0f7d1fc60 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDaoTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/repository/dao/JdbcJobExecutionDaoTests.java @@ -1,8 +1,13 @@ package org.springframework.batch.core.repository.dao; +import static org.junit.Assert.*; + +import org.springframework.batch.core.BatchStatus; 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.Test; import org.junit.runner.RunWith; @RunWith(SpringJUnit4ClassRunner.class) @@ -32,5 +37,5 @@ public class JdbcJobExecutionDaoTests extends AbstractJobExecutionDaoTests { protected StepExecutionDao getStepExecutionDao() { return stepExecutionDao; } - + } diff --git a/spring-batch-core/src/test/java/org/springframework/batch/core/repository/support/SimpleJobRepositoryTests.java b/spring-batch-core/src/test/java/org/springframework/batch/core/repository/support/SimpleJobRepositoryTests.java index f99ae6492..e92d2ae60 100644 --- a/spring-batch-core/src/test/java/org/springframework/batch/core/repository/support/SimpleJobRepositoryTests.java +++ b/spring-batch-core/src/test/java/org/springframework/batch/core/repository/support/SimpleJobRepositoryTests.java @@ -24,6 +24,7 @@ import java.util.List; import org.junit.Before; import org.junit.Test; +import org.springframework.batch.core.BatchStatus; import org.springframework.batch.core.JobExecution; import org.springframework.batch.core.JobInstance; import org.springframework.batch.core.JobParameters; @@ -184,5 +185,16 @@ public class SimpleJobRepositoryTests { long lastUpdated = stepExecution.getLastUpdated().getTime(); assertTrue(lastUpdated > (before - 1000)); } + + @Test + public void testInterrupted(){ + + jobExecution.setStatus(BatchStatus.STOPPING); + StepExecution stepExecution = new StepExecution("stepName", jobExecution); + stepExecution.setId(323L); + + jobRepository.update(stepExecution); + assertTrue(stepExecution.isTerminateOnly()); + } }