RESOLVED - issue BATCH-1197: Rerunning a job sometimes creates new job instance

This commit is contained in:
dsyer
2009-04-07 09:05:14 +00:00
parent 46a912cd6e
commit f12edbb546
2 changed files with 143 additions and 90 deletions

View File

@@ -25,6 +25,7 @@ import java.sql.SQLException;
import java.sql.Timestamp;
import java.sql.Types;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -58,7 +59,8 @@ import org.springframework.util.StringUtils;
* @author Dave Syer
* @author Robert Kasanicky
*/
public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements JobInstanceDao, InitializingBean {
public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
JobInstanceDao, InitializingBean {
private static final String CREATE_JOB_INSTANCE = "INSERT into %PREFIX%JOB_INSTANCE(JOB_INSTANCE_ID, JOB_NAME, JOB_KEY, VERSION)"
+ " values (?, ?, ?, ?)";
@@ -68,7 +70,8 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
private static final String FIND_JOBS_WITH_NAME = "SELECT JOB_INSTANCE_ID, JOB_NAME from %PREFIX%JOB_INSTANCE where JOB_NAME = ?";
private static final String FIND_JOBS_WITH_KEY = FIND_JOBS_WITH_NAME + " and JOB_KEY = ?";
private static final String FIND_JOBS_WITH_KEY = FIND_JOBS_WITH_NAME
+ " and JOB_KEY = ?";
private static final String FIND_JOBS_WITH_EMPTY_KEY = "SELECT JOB_INSTANCE_ID, JOB_NAME from %PREFIX%JOB_INSTANCE where JOB_NAME = ? and (JOB_KEY = ? OR JOB_KEY is NULL)";
@@ -92,24 +95,30 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
* then passing the Id and parameter values into an INSERT statement.
*
* @see JobInstanceDao#createJobInstance(String, JobParameters)
* @throws IllegalArgumentException if any {@link JobParameters} fields are
* null.
* @throws IllegalArgumentException
* if any {@link JobParameters} fields are null.
*/
public JobInstance createJobInstance(String jobName, JobParameters jobParameters) {
public JobInstance createJobInstance(String jobName,
JobParameters jobParameters) {
Assert.notNull(jobName, "Job name must not be null.");
Assert.notNull(jobParameters, "JobParameters must not be null.");
Assert.state(getJobInstance(jobName, jobParameters) == null, "JobInstance must not already exist");
Assert.state(getJobInstance(jobName, jobParameters) == null,
"JobInstance must not already exist");
Long jobId = jobIncrementer.nextLongValue();
JobInstance jobInstance = new JobInstance(jobId, jobParameters, jobName);
jobInstance.incrementVersion();
Object[] parameters = new Object[] { jobId, jobName, createJobKey(jobParameters), jobInstance.getVersion() };
getJdbcTemplate().getJdbcOperations().update(getQuery(CREATE_JOB_INSTANCE), parameters,
new int[] { Types.INTEGER, Types.VARCHAR, Types.VARCHAR, Types.INTEGER });
Object[] parameters = new Object[] { jobId, jobName,
createJobKey(jobParameters), jobInstance.getVersion() };
getJdbcTemplate().getJdbcOperations().update(
getQuery(CREATE_JOB_INSTANCE),
parameters,
new int[] { Types.INTEGER, Types.VARCHAR, Types.VARCHAR,
Types.INTEGER });
insertJobParameters(jobId, jobParameters);
@@ -120,24 +129,27 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
Map<String, JobParameter> props = jobParameters.getParameters();
StringBuffer stringBuffer = new StringBuffer();
for (Entry<String, JobParameter> entry : props.entrySet()) {
stringBuffer.append(entry.toString() + ";");
List<String> keys = new ArrayList<String>(props.keySet());
Collections.sort(keys);
for (String key : keys) {
stringBuffer.append(key + "=" + props.get(key).toString() + ";");
}
MessageDigest digest;
try {
digest = MessageDigest.getInstance("MD5");
} catch (NoSuchAlgorithmException e) {
throw new IllegalStateException(
"MD5 algorithm not available. Fatal (should be in the JDK).");
}
catch (NoSuchAlgorithmException e) {
throw new IllegalStateException("MD5 algorithm not available. Fatal (should be in the JDK).");
}
try {
byte[] bytes = digest.digest(stringBuffer.toString().getBytes("UTF-8"));
byte[] bytes = digest.digest(stringBuffer.toString().getBytes(
"UTF-8"));
return String.format("%032x", new BigInteger(1, bytes));
}
catch (UnsupportedEncodingException e) {
throw new IllegalStateException("UTF-8 encoding not available. Fatal (should be in the JDK).");
} catch (UnsupportedEncodingException e) {
throw new IllegalStateException(
"UTF-8 encoding not available. Fatal (should be in the JDK).");
}
}
@@ -148,9 +160,11 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
*/
private void insertJobParameters(Long jobId, JobParameters jobParameters) {
for (Entry<String, JobParameter> entry : jobParameters.getParameters().entrySet()) {
for (Entry<String, JobParameter> entry : jobParameters.getParameters()
.entrySet()) {
JobParameter jobParameter = entry.getValue();
insertParameter(jobId, jobParameter.getType(), entry.getKey(), jobParameter.getValue());
insertParameter(jobId, jobParameter.getType(), entry.getKey(),
jobParameter.getValue());
}
}
@@ -158,26 +172,29 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
* Convenience method that inserts an individual records into the
* JobParameters table.
*/
private void insertParameter(Long jobId, ParameterType type, String key, Object value) {
private void insertParameter(Long jobId, ParameterType type, String key,
Object value) {
Object[] args = new Object[0];
int[] argTypes = new int[] { Types.BIGINT, Types.VARCHAR, Types.VARCHAR, Types.VARCHAR, Types.TIMESTAMP,
Types.BIGINT, Types.DOUBLE };
int[] argTypes = new int[] { Types.BIGINT, Types.VARCHAR,
Types.VARCHAR, Types.VARCHAR, Types.TIMESTAMP, Types.BIGINT,
Types.DOUBLE };
if (type == ParameterType.STRING) {
args = new Object[] { jobId, key, type, value, new Timestamp(0L), 0L, 0D };
}
else if (type == ParameterType.LONG) {
args = new Object[] { jobId, key, type, "", new Timestamp(0L), value, new Double(0) };
}
else if (type == ParameterType.DOUBLE) {
args = new Object[] { jobId, key, type, "", new Timestamp(0L), 0L, value };
}
else if (type == ParameterType.DATE) {
args = new Object[] { jobId, key, type, value, new Timestamp(0L),
0L, 0D };
} else if (type == ParameterType.LONG) {
args = new Object[] { jobId, key, type, "", new Timestamp(0L),
value, new Double(0) };
} else if (type == ParameterType.DOUBLE) {
args = new Object[] { jobId, key, type, "", new Timestamp(0L), 0L,
value };
} else if (type == ParameterType.DATE) {
args = new Object[] { jobId, key, type, "", value, 0L, 0D };
}
getJdbcTemplate().getJdbcOperations().update(getQuery(CREATE_JOB_PARAMETERS), args, argTypes);
getJdbcTemplate().getJdbcOperations().update(
getQuery(CREATE_JOB_PARAMETERS), args, argTypes);
}
/**
@@ -185,30 +202,33 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
* given identifier, adding them to a list via the RowMapper callback.
*
* @see JobInstanceDao#getJobInstance(String, JobParameters)
* @throws IllegalArgumentException if any {@link JobParameters} fields are
* null.
* @throws IllegalArgumentException
* if any {@link JobParameters} fields are null.
*/
public JobInstance getJobInstance(final String jobName, final JobParameters jobParameters) {
public JobInstance getJobInstance(final String jobName,
final JobParameters jobParameters) {
Assert.notNull(jobName, "Job name must not be null.");
Assert.notNull(jobParameters, "JobParameters must not be null.");
String jobKey = createJobKey(jobParameters);
ParameterizedRowMapper<JobInstance> rowMapper = new JobInstanceRowMapper(jobParameters);
ParameterizedRowMapper<JobInstance> rowMapper = new JobInstanceRowMapper(
jobParameters);
List<JobInstance> instances;
if (StringUtils.hasLength(jobKey)) {
instances = getJdbcTemplate().query(getQuery(FIND_JOBS_WITH_KEY), rowMapper, jobName, jobKey);
}
else {
instances = getJdbcTemplate().query(getQuery(FIND_JOBS_WITH_EMPTY_KEY), rowMapper, jobName, jobKey);
instances = getJdbcTemplate().query(getQuery(FIND_JOBS_WITH_KEY),
rowMapper, jobName, jobKey);
} else {
instances = getJdbcTemplate().query(
getQuery(FIND_JOBS_WITH_EMPTY_KEY), rowMapper, jobName,
jobKey);
}
if (instances.isEmpty()) {
return null;
}
else {
} else {
Assert.state(instances.size() == 1);
return instances.get(0);
}
@@ -224,9 +244,9 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
public JobInstance getJobInstance(Long instanceId) {
try {
return getJdbcTemplate().queryForObject(getQuery(GET_JOB_FROM_ID), new JobInstanceRowMapper(), instanceId);
}
catch (EmptyResultDataAccessException e) {
return getJdbcTemplate().queryForObject(getQuery(GET_JOB_FROM_ID),
new JobInstanceRowMapper(), instanceId);
} catch (EmptyResultDataAccessException e) {
return null;
}
@@ -244,22 +264,20 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
JobParameter value = null;
if (type == ParameterType.STRING) {
value = new JobParameter(rs.getString(4));
}
else if (type == ParameterType.LONG) {
} else if (type == ParameterType.LONG) {
value = new JobParameter(rs.getLong(6));
}
else if (type == ParameterType.DOUBLE) {
} else if (type == ParameterType.DOUBLE) {
value = new JobParameter(rs.getDouble(7));
}
else if (type == ParameterType.DATE) {
} else if (type == ParameterType.DATE) {
value = new JobParameter(rs.getTimestamp(5));
}
// No need to assert that value is not null because it's an enum
map.put(rs.getString(2), value);
}
};
getJdbcTemplate().getJdbcOperations()
.query(getQuery(FIND_PARAMS_FROM_ID), new Object[] { instanceId }, handler);
getJdbcTemplate().getJdbcOperations().query(
getQuery(FIND_PARAMS_FROM_ID), new Object[] { instanceId },
handler);
return new JobParameters(map);
}
@@ -271,11 +289,13 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
* ()
*/
public List<String> getJobNames() {
return getJdbcTemplate().query(getQuery(FIND_JOB_NAMES), new ParameterizedRowMapper<String>() {
public String mapRow(ResultSet rs, int rowNum) throws SQLException {
return rs.getString(1);
}
});
return getJdbcTemplate().query(getQuery(FIND_JOB_NAMES),
new ParameterizedRowMapper<String>() {
public String mapRow(ResultSet rs, int rowNum)
throws SQLException {
return rs.getString(1);
}
});
}
/*
@@ -284,13 +304,15 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
* @seeorg.springframework.batch.core.repository.dao.JobInstanceDao#
* getLastJobInstances(java.lang.String, int)
*/
public List<JobInstance> getJobInstances(String jobName, final int start, final int count) {
public List<JobInstance> getJobInstances(String jobName, final int start,
final int count) {
ResultSetExtractor extractor = new ResultSetExtractor() {
private List<JobInstance> list = new ArrayList<JobInstance>();
public Object extractData(ResultSet rs) throws SQLException, DataAccessException {
public Object extractData(ResultSet rs) throws SQLException,
DataAccessException {
int rowNum = 0;
while (rowNum < start && rs.next()) {
rowNum++;
@@ -306,8 +328,9 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
};
@SuppressWarnings("unchecked")
List<JobInstance> result = (List<JobInstance>) getJdbcTemplate().getJdbcOperations().query(
getQuery(FIND_LAST_JOBS_BY_NAME), new Object[] { jobName }, extractor);
List<JobInstance> result = (List<JobInstance>) getJdbcTemplate()
.getJdbcOperations().query(getQuery(FIND_LAST_JOBS_BY_NAME),
new Object[] { jobName }, extractor);
return result;
}
@@ -322,10 +345,10 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
public JobInstance getJobInstance(JobExecution jobExecution) {
try {
return getJdbcTemplate().queryForObject(getQuery(GET_JOB_FROM_EXECUTION_ID), new JobInstanceRowMapper(),
jobExecution.getId());
}
catch (EmptyResultDataAccessException e) {
return getJdbcTemplate().queryForObject(
getQuery(GET_JOB_FROM_EXECUTION_ID),
new JobInstanceRowMapper(), jobExecution.getId());
} catch (EmptyResultDataAccessException e) {
return null;
}
}
@@ -334,7 +357,8 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
* Setter for {@link DataFieldMaxValueIncrementer} to be used when
* generating primary keys for {@link JobInstance} instances.
*
* @param jobIncrementer the {@link DataFieldMaxValueIncrementer}
* @param jobIncrementer
* the {@link DataFieldMaxValueIncrementer}
*/
public void setJobIncrementer(DataFieldMaxValueIncrementer jobIncrementer) {
this.jobIncrementer = jobIncrementer;
@@ -349,7 +373,8 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
* @author Dave Syer
*
*/
private final class JobInstanceRowMapper implements ParameterizedRowMapper<JobInstance> {
private final class JobInstanceRowMapper implements
ParameterizedRowMapper<JobInstance> {
private JobParameters jobParameters;
@@ -365,7 +390,8 @@ public class JdbcJobInstanceDao extends AbstractJdbcBatchMetadataDao implements
if (jobParameters == null) {
jobParameters = getJobParameters(id);
}
JobInstance jobInstance = new JobInstance(rs.getLong(1), jobParameters, rs.getString(2));
JobInstance jobInstance = new JobInstance(rs.getLong(1),
jobParameters, rs.getString(2));
// should always be at version=0 because they never get updated
jobInstance.incrementVersion();
return jobInstance;

View File

@@ -1,43 +1,70 @@
package org.springframework.batch.core.repository.dao;
import static org.junit.Assert.*;
import static org.junit.Assert.assertEquals;
import java.math.BigInteger;
import java.security.MessageDigest;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.batch.core.JobExecution;
import org.springframework.batch.core.JobInstance;
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.junit.Assert;
import org.junit.Test;
import org.junit.runner.RunWith;
@RunWith(SpringJUnit4ClassRunner.class)
@ContextConfiguration(locations = "sql-dao-test.xml")
public class JdbcJobInstanceDaoTests extends AbstractJobInstanceDaoTests {
protected JobInstanceDao getJobInstanceDao() {
deleteFromTables("BATCH_JOB_EXECUTION_CONTEXT", "BATCH_STEP_EXECUTION_CONTEXT", "BATCH_STEP_EXECUTION", "BATCH_JOB_EXECUTION",
"BATCH_JOB_PARAMS", "BATCH_JOB_INSTANCE");
deleteFromTables("BATCH_JOB_EXECUTION_CONTEXT",
"BATCH_STEP_EXECUTION_CONTEXT", "BATCH_STEP_EXECUTION",
"BATCH_JOB_EXECUTION", "BATCH_JOB_PARAMS", "BATCH_JOB_INSTANCE");
return (JobInstanceDao) applicationContext.getBean("jobInstanceDao");
}
@Test
public void testFindJobInstanceByExecution(){
JobExecutionDao jobExecutionDao = (JobExecutionDao) applicationContext.getBean("jobExecutionDao");
JobInstance jobInstance = dao.createJobInstance("testInstance", new JobParameters());
public void testFindJobInstanceByExecution() {
JobExecutionDao jobExecutionDao = (JobExecutionDao) applicationContext
.getBean("jobExecutionDao");
JobInstance jobInstance = dao.createJobInstance("testInstance",
new JobParameters());
JobExecution jobExecution = new JobExecution(jobInstance, 2L);
jobExecutionDao.saveJobExecution(jobExecution);
JobInstance returnedInstance = dao.getJobInstance(jobExecution);
Assert.assertEquals(jobInstance, returnedInstance);
assertEquals(jobInstance, returnedInstance);
}
@Test
public void testCreateJobKey() {
JdbcJobInstanceDao jdbcDao = (JdbcJobInstanceDao) dao;
JobParameters jobParameters = new JobParametersBuilder().addString(
"foo", "bar").addString("bar", "foo").toJobParameters();
String key = jdbcDao.createJobKey(jobParameters);
assertEquals(32, key.length());
}
@Test
public void testCreateJobKeyOrdering() {
JdbcJobInstanceDao jdbcDao = (JdbcJobInstanceDao) dao;
JobParameters jobParameters1 = new JobParametersBuilder().addString(
"foo", "bar").addString("bar", "foo").toJobParameters();
String key1 = jdbcDao.createJobKey(jobParameters1);
JobParameters jobParameters2 = new JobParametersBuilder().addString(
"bar", "foo").addString("foo", "bar").toJobParameters();
String key2 = jdbcDao.createJobKey(jobParameters2);
assertEquals(key1, key2);
}
@Test
public void testHexing() throws Exception {
MessageDigest digest = MessageDigest.getInstance("MD5");
@@ -46,9 +73,9 @@ public class JdbcJobInstanceDaoTests extends AbstractJobInstanceDaoTests {
for (byte bite : bytes) {
output.append(String.format("%02x", bite));
}
assertEquals("Wrong hash: "+output,32, output.length());
assertEquals("Wrong hash: " + output, 32, output.length());
String value = String.format("%032x", new BigInteger(1, bytes));
assertEquals("Wrong hash: "+value,32, value.length());
assertEquals("Wrong hash: " + value, 32, value.length());
assertEquals(value, output.toString());
}
}