BATCH-401: prototype of JobInterruptionAdvice added to samples (see GracefulShutdownTests)
This commit is contained in:
@@ -0,0 +1,44 @@
|
||||
package org.springframework.batch.sample.advice;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.springframework.batch.core.domain.BatchStatus;
|
||||
import org.springframework.batch.core.domain.JobExecution;
|
||||
import org.springframework.batch.core.domain.JobInterruptedException;
|
||||
|
||||
/**
|
||||
* Monitors {@link JobExecution} and throws a {@link JobInterruptedException} in
|
||||
* case the execution has been requested to stop using
|
||||
* {@link JobExecution#stop()} or executing thread has been interrupted.
|
||||
*
|
||||
* @author Robert Kasanicky
|
||||
*/
|
||||
public class DefaultJobInterruptionAdvice implements JobInterruptionAdvice {
|
||||
|
||||
private static final Log logger = LogFactory.getLog(DefaultJobInterruptionAdvice.class);
|
||||
|
||||
private JobExecution jobExecution;
|
||||
|
||||
/*
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.batch.core.runtime.JobInterruptionAdvice#checkInterrupt()
|
||||
*/
|
||||
public void checkInterrupt() throws JobInterruptedException {
|
||||
|
||||
logger.info("checking job execution for interrupt");
|
||||
|
||||
if (Thread.currentThread().isInterrupted() || jobExecution.getStatus() == BatchStatus.STOPPING) {
|
||||
throw new JobInterruptedException("Job execution interrupted by user");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Setter for JobExecution - to be applied as 'after returning' advice that
|
||||
* captures return value on the method that creates {@link JobExecution} for
|
||||
* the job run.
|
||||
*/
|
||||
public void setJobExecution(JobExecution jobExecution) {
|
||||
this.jobExecution = jobExecution;
|
||||
logger.info("JobExecution set");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,25 @@
|
||||
package org.springframework.batch.sample.advice;
|
||||
|
||||
import org.springframework.batch.core.domain.JobInterruptedException;
|
||||
|
||||
/**
|
||||
* Interface for monitors that check whether a job was interrupted by user.
|
||||
*
|
||||
* Interruption checking is a cross-cutting concern, therefore favorably handled
|
||||
* by AOP. The implementation if expected to be used as 'before advice' or
|
||||
* 'after advice' or both.
|
||||
*
|
||||
* @see DefaultJobInterruptionAdvice
|
||||
*
|
||||
* @author Robert Kasanicky
|
||||
*/
|
||||
public interface JobInterruptionAdvice {
|
||||
|
||||
/**
|
||||
* No-op unless job was interrupted by user.
|
||||
*
|
||||
* @throws JobInterruptedException if job was interrupted
|
||||
*/
|
||||
void checkInterrupt() throws JobInterruptedException;
|
||||
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
package org.springframework.batch.sample.item.reader;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.exception.MarkFailedException;
|
||||
import org.springframework.batch.item.exception.ResetFailedException;
|
||||
|
||||
|
||||
/**
|
||||
* Creates items infinitely - useful for testing job interruption by user.
|
||||
*/
|
||||
public class InfiniteItemReader implements ItemReader {
|
||||
|
||||
private static final Log logger = LogFactory.getLog(InfiniteItemReader.class);
|
||||
|
||||
private long counter = 0;
|
||||
|
||||
public Object read() throws Exception {
|
||||
Object item = "item" + counter;
|
||||
Thread.sleep(10);
|
||||
logger.info("read item: " + item);
|
||||
return item;
|
||||
}
|
||||
|
||||
public void mark() throws MarkFailedException {
|
||||
// no-op
|
||||
|
||||
}
|
||||
|
||||
public void reset() throws ResetFailedException {
|
||||
// no-op
|
||||
|
||||
}
|
||||
|
||||
|
||||
|
||||
}
|
||||
@@ -0,0 +1,15 @@
|
||||
package org.springframework.batch.sample.step.support;
|
||||
|
||||
import org.springframework.batch.core.domain.JobInterruptedException;
|
||||
import org.springframework.batch.execution.step.support.StepInterruptionPolicy;
|
||||
import org.springframework.batch.repeat.RepeatContext;
|
||||
|
||||
public class NoopStepInterruptionPolicy implements StepInterruptionPolicy {
|
||||
|
||||
public void checkInterrupted(RepeatContext context)
|
||||
throws JobInterruptedException {
|
||||
// no-op
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
@@ -20,12 +20,41 @@
|
||||
|
||||
<bean id="infiniteLoopJob" parent="simpleJob">
|
||||
<property name="steps">
|
||||
<bean id="step1" parent="taskletStep">
|
||||
<!--bean id="step1" parent="taskletStep">
|
||||
<property name="tasklet">
|
||||
<bean id="module" class="org.springframework.batch.sample.tasklet.InfiniteLoopTasklet"/>
|
||||
</property>
|
||||
</bean-->
|
||||
<bean id="infiniteStep" parent="simpleStep">
|
||||
<property name="chunkOperations">
|
||||
<bean class="org.springframework.batch.repeat.support.RepeatTemplate" />
|
||||
</property>
|
||||
<property name="interruptionPolicy">
|
||||
<bean class="org.springframework.batch.sample.step.support.NoopStepInterruptionPolicy" />
|
||||
</property>
|
||||
<property name="itemReader">
|
||||
<bean id="reader" class="org.springframework.batch.sample.item.reader.InfiniteItemReader" />
|
||||
</property>
|
||||
<property name="itemWriter">
|
||||
<bean id="writer" class="org.springframework.batch.sample.item.writer.DummyItemWriter" />
|
||||
</property>
|
||||
</bean>
|
||||
</property>
|
||||
</bean>
|
||||
|
||||
<bean id="jobInterruptAdvice" class="org.springframework.batch.sample.advice.DefaultJobInterruptionAdvice" />
|
||||
|
||||
<aop:config>
|
||||
<aop:aspect id="interruptChecker" ref="jobInterruptAdvice" >
|
||||
|
||||
<aop:after-returning method="setJobExecution"
|
||||
pointcut="execution(* org.springframework.batch.core.repository.JobRepository.createJobExecution(..))"
|
||||
returning="jobExecution" />
|
||||
|
||||
<aop:before method="checkInterrupt"
|
||||
pointcut="execution(* org.springframework.batch.repeat.RepeatOperations.iterate(..))"/>
|
||||
|
||||
</aop:aspect>
|
||||
</aop:config>
|
||||
|
||||
</beans>
|
||||
@@ -22,7 +22,7 @@ import org.springframework.batch.core.domain.JobParameters;
|
||||
|
||||
/**
|
||||
* Functional test for graceful shutdown. A batch container is started in a new thread,
|
||||
* then it's stopped via the Lifecycle interface.
|
||||
* then it's stopped using {@link JobExecution#stop()}.
|
||||
*
|
||||
* @author Lucas Ward
|
||||
*
|
||||
@@ -39,7 +39,7 @@ public class GracefulShutdownFunctionalTests extends AbstractBatchLauncherTests
|
||||
|
||||
JobExecution jobExecution = launcher.run(getJob(), jobParameters);
|
||||
|
||||
Thread.sleep(500);
|
||||
Thread.sleep(1000);
|
||||
|
||||
assertEquals(BatchStatus.STARTED, jobExecution.getStatus());
|
||||
assertTrue(jobExecution.isRunning());
|
||||
@@ -49,16 +49,11 @@ public class GracefulShutdownFunctionalTests extends AbstractBatchLauncherTests
|
||||
int count = 0;
|
||||
while(jobExecution.isRunning() && count <= 10){
|
||||
logger.info("Checking for end time in JobExecution: count="+count);
|
||||
Thread.sleep(10);
|
||||
Thread.sleep(100);
|
||||
count++;
|
||||
}
|
||||
if (count>10) {
|
||||
// TODO: fix this
|
||||
// fail("Timed out waiting for job to end.");
|
||||
}
|
||||
|
||||
// TODO: fix this
|
||||
// assertFalse(jobExecution.isRunning());
|
||||
assertFalse("Timed out waiting for job to end.", jobExecution.isRunning());
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,92 @@
|
||||
package org.springframework.batch.sample.advice;
|
||||
|
||||
import junit.framework.TestCase;
|
||||
|
||||
import org.springframework.batch.core.domain.BatchStatus;
|
||||
import org.springframework.batch.core.domain.JobExecution;
|
||||
import org.springframework.batch.core.domain.JobInstance;
|
||||
import org.springframework.batch.core.domain.JobInterruptedException;
|
||||
import org.springframework.batch.core.domain.JobParameters;
|
||||
import org.springframework.batch.sample.tasklet.JobSupport;
|
||||
|
||||
/**
|
||||
* Tests for {@link DefaultJobInterruptionAdvice}.
|
||||
*
|
||||
* @author Robert Kasanicky
|
||||
*/
|
||||
public class DefaultJobInterruptionAdviceTests extends TestCase {
|
||||
|
||||
private DefaultJobInterruptionAdvice tested = new DefaultJobInterruptionAdvice();
|
||||
|
||||
private JobExecution jobExecution = new JobExecution(new JobInstance(new Long(1), new JobParameters(),
|
||||
new JobSupport("interruptJob")));
|
||||
|
||||
protected void setUp() throws Exception {
|
||||
tested.setJobExecution(jobExecution);
|
||||
}
|
||||
|
||||
/**
|
||||
* Scenario when JobExecution is requested to stop explicitly.
|
||||
*/
|
||||
public void testStop() throws JobInterruptedException {
|
||||
|
||||
jobExecution.stop();
|
||||
try {
|
||||
tested.checkInterrupt();
|
||||
fail();
|
||||
}
|
||||
catch (JobInterruptedException e) {
|
||||
// expected
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Scenario when executing thread is interrupted. Needs to be run in
|
||||
* separate victim thread - interrupting current thread would affect other
|
||||
* tests.
|
||||
*/
|
||||
public void testThreadInterrupt() throws Exception {
|
||||
|
||||
class TestJob implements Runnable {
|
||||
|
||||
volatile boolean interrupted = false;
|
||||
|
||||
public void run() {
|
||||
while (true) {
|
||||
try {
|
||||
tested.checkInterrupt();
|
||||
}
|
||||
catch (JobInterruptedException expected) {
|
||||
interrupted = true;
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
TestJob job = new TestJob();
|
||||
Thread victim = new Thread(job);
|
||||
|
||||
victim.start();
|
||||
victim.interrupt();
|
||||
|
||||
Thread.sleep(1000);
|
||||
assertTrue(job.interrupted);
|
||||
|
||||
}
|
||||
|
||||
/**
|
||||
* No exception raised when JobExecution has non-terminating status.
|
||||
*/
|
||||
public void testNoInterrupt() throws JobInterruptedException {
|
||||
|
||||
BatchStatus[] notInterruptedValues = { BatchStatus.STARTED, BatchStatus.STARTING, BatchStatus.STOPPED,
|
||||
BatchStatus.UNKNOWN, BatchStatus.COMPLETED, BatchStatus.FAILED };
|
||||
|
||||
for (int i = 0; i < notInterruptedValues.length; i++) {
|
||||
jobExecution.setStatus(notInterruptedValues[i]);
|
||||
tested.checkInterrupt();
|
||||
assertTrue(true);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user