PollableSource died (apparently)
This commit is contained in:
@@ -2,7 +2,6 @@ package org.springframework.batch.integration.chunk;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
import java.util.Arrays;
|
||||
|
||||
@@ -15,7 +14,6 @@ import org.springframework.batch.core.JobExecution;
|
||||
import org.springframework.batch.core.JobParametersBuilder;
|
||||
import org.springframework.batch.core.Step;
|
||||
import org.springframework.batch.core.StepExecution;
|
||||
import org.springframework.batch.core.UnexpectedJobExecutionException;
|
||||
import org.springframework.batch.core.job.SimpleJob;
|
||||
import org.springframework.batch.core.repository.JobExecutionAlreadyRunningException;
|
||||
import org.springframework.batch.core.repository.JobInstanceAlreadyCompleteException;
|
||||
@@ -54,7 +52,7 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
@Qualifier("replies")
|
||||
private PollableChannel replies;
|
||||
|
||||
private SimpleStepFactoryBean<Object,Object> factory = new SimpleStepFactoryBean<Object,Object>();
|
||||
private SimpleStepFactoryBean<Object, Object> factory = new SimpleStepFactoryBean<Object, Object>();
|
||||
|
||||
private SimpleJobRepository jobRepository;
|
||||
|
||||
@@ -63,8 +61,8 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
@Before
|
||||
public void setUp() {
|
||||
|
||||
jobRepository = new SimpleJobRepository(new MapJobInstanceDao(),
|
||||
new MapJobExecutionDao(), new MapStepExecutionDao(), new MapExecutionContextDao());
|
||||
jobRepository = new SimpleJobRepository(new MapJobInstanceDao(), new MapJobExecutionDao(),
|
||||
new MapStepExecutionDao(), new MapExecutionContextDao());
|
||||
factory.setJobRepository(jobRepository);
|
||||
factory.setTransactionManager(new ResourcelessTransactionManager());
|
||||
factory.setBeanName("step");
|
||||
@@ -75,10 +73,10 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
writer.setOutputChannel(requests);
|
||||
|
||||
TestItemWriter.count = 0;
|
||||
|
||||
|
||||
// Drain queues
|
||||
Message<?> message = replies.receive(10);
|
||||
while (message!=null) {
|
||||
while (message != null) {
|
||||
System.err.println(message);
|
||||
message = replies.receive(10);
|
||||
}
|
||||
@@ -87,7 +85,8 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
|
||||
@After
|
||||
public void tearDown() {
|
||||
while(replies.receive(10L)!=null) {}
|
||||
while (replies.receive(10L) != null) {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -100,10 +99,8 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
ExecutionContext executionContext = new ExecutionContext();
|
||||
writer.update(executionContext);
|
||||
writer.open(executionContext);
|
||||
assertEquals(0, executionContext
|
||||
.getLong(ChunkMessageChannelItemWriter.EXPECTED));
|
||||
assertEquals(0, executionContext
|
||||
.getLong(ChunkMessageChannelItemWriter.ACTUAL));
|
||||
assertEquals(0, executionContext.getLong(ChunkMessageChannelItemWriter.EXPECTED));
|
||||
assertEquals(0, executionContext.getLong(ChunkMessageChannelItemWriter.ACTUAL));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -116,7 +113,7 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
|
||||
StepExecution stepExecution = getStepExecution(step);
|
||||
step.execute(stepExecution);
|
||||
|
||||
|
||||
waitForResults(6, 10);
|
||||
|
||||
assertEquals(6, TestItemWriter.count);
|
||||
@@ -135,10 +132,8 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
StepExecution stepExecution = getStepExecution(step);
|
||||
|
||||
// Set up context with two messages (chunks) in the backlog
|
||||
stepExecution.getExecutionContext().putLong(
|
||||
ChunkMessageChannelItemWriter.EXPECTED, 6);
|
||||
stepExecution.getExecutionContext().putLong(
|
||||
ChunkMessageChannelItemWriter.ACTUAL, 4);
|
||||
stepExecution.getExecutionContext().putLong(ChunkMessageChannelItemWriter.EXPECTED, 6);
|
||||
stepExecution.getExecutionContext().putLong(ChunkMessageChannelItemWriter.ACTUAL, 4);
|
||||
// And make the back log real
|
||||
requests.send(getSimpleMessage("foo", stepExecution.getJobExecution().getJobId()));
|
||||
requests.send(getSimpleMessage("bar", stepExecution.getJobExecution().getJobId()));
|
||||
@@ -162,19 +157,15 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
StepExecution stepExecution = getStepExecution(step);
|
||||
|
||||
// Set up context with two messages (chunks) in the backlog
|
||||
stepExecution.getExecutionContext().putLong(
|
||||
ChunkMessageChannelItemWriter.EXPECTED, 3);
|
||||
stepExecution.getExecutionContext().putLong(
|
||||
ChunkMessageChannelItemWriter.ACTUAL, 2);
|
||||
stepExecution.getExecutionContext().putLong(ChunkMessageChannelItemWriter.EXPECTED, 3);
|
||||
stepExecution.getExecutionContext().putLong(ChunkMessageChannelItemWriter.ACTUAL, 2);
|
||||
// And make the back log real
|
||||
requests.send(getSimpleMessage("foo", new Long(4321)));
|
||||
try {
|
||||
step.execute(stepExecution);
|
||||
fail("Expected UnexpectedJobExecutionException");
|
||||
} catch (UnexpectedJobExecutionException e) {
|
||||
String message = e.getCause().getMessage();
|
||||
assertTrue("Message does not contain 'wrong job': "+message, message.contains("wrong job"));
|
||||
}
|
||||
step.execute(stepExecution);
|
||||
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
|
||||
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
|
||||
String message = stepExecution.getExitStatus().getExitDescription();
|
||||
assertTrue("Message does not contain 'wrong job': " + message, message.contains("wrong job"));
|
||||
|
||||
waitForResults(1, 10);
|
||||
|
||||
@@ -184,14 +175,13 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
}
|
||||
|
||||
/**
|
||||
* @param jobId
|
||||
* @param string
|
||||
* @param jobId
|
||||
* @param string
|
||||
* @return
|
||||
*/
|
||||
@SuppressWarnings("unchecked")
|
||||
private GenericMessage<ChunkRequest> getSimpleMessage(String string, Long jobId) {
|
||||
ChunkRequest chunk = new ChunkRequest(StringUtils
|
||||
.commaDelimitedListToSet(string), jobId, 0);
|
||||
ChunkRequest chunk = new ChunkRequest(StringUtils.commaDelimitedListToSet(string), jobId, 0);
|
||||
GenericMessage<ChunkRequest> message = new GenericMessage<ChunkRequest>(chunk);
|
||||
return message;
|
||||
}
|
||||
@@ -206,12 +196,11 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
Step step = (Step) factory.getObject();
|
||||
|
||||
StepExecution stepExecution = getStepExecution(step);
|
||||
try {
|
||||
step.execute(stepExecution);
|
||||
fail("Expected AsynchronousFailureException");
|
||||
} catch (AsynchronousFailureException e) {
|
||||
assertTrue(e.getMessage().contains("bad"));
|
||||
}
|
||||
step.execute(stepExecution);
|
||||
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
|
||||
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
|
||||
String message = stepExecution.getExitStatus().getExitDescription();
|
||||
assertTrue("Message does not contain 'bad': " + message, message.contains("bad"));
|
||||
|
||||
waitForResults(2, 10);
|
||||
|
||||
@@ -235,23 +224,18 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
StepExecution stepExecution = getStepExecution(step);
|
||||
|
||||
// Set up expectation of three messages (chunks) in the backlog
|
||||
stepExecution.getExecutionContext().putLong(
|
||||
ChunkMessageChannelItemWriter.EXPECTED, 6);
|
||||
stepExecution.getExecutionContext().putLong(
|
||||
ChunkMessageChannelItemWriter.ACTUAL, 3);
|
||||
stepExecution.getExecutionContext().putLong(ChunkMessageChannelItemWriter.EXPECTED, 6);
|
||||
stepExecution.getExecutionContext().putLong(ChunkMessageChannelItemWriter.ACTUAL, 3);
|
||||
/*
|
||||
* With no backlog we process all the items, but the listener can't
|
||||
* reconcile the expected number of items with the actual. An infinite
|
||||
* loop would be bad, so the best we can do is fail as fast as possible.
|
||||
*/
|
||||
try {
|
||||
step.execute(stepExecution);
|
||||
fail("Expected UnexpectedJobExecutionException");
|
||||
} catch (UnexpectedJobExecutionException e) {
|
||||
String message = e.getCause().getMessage();
|
||||
assertTrue("Message did not contain 'timed out': " + message,
|
||||
message.toLowerCase().contains("timed out"));
|
||||
}
|
||||
step.execute(stepExecution);
|
||||
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
|
||||
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
|
||||
String message = stepExecution.getExitStatus().getExitDescription();
|
||||
assertTrue("Message did not contain 'timed out': " + message, message.toLowerCase().contains("timed out"));
|
||||
|
||||
assertEquals(0, TestItemWriter.count);
|
||||
assertEquals(0, stepExecution.getReadCount());
|
||||
@@ -283,13 +267,10 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
assertTrue(6 >= TestItemWriter.count);
|
||||
|
||||
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
|
||||
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution
|
||||
.getExitStatus().getExitCode());
|
||||
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
|
||||
|
||||
String exitDescription = stepExecution.getExitStatus()
|
||||
.getExitDescription();
|
||||
assertTrue("Exit description does not contain exception type name: "
|
||||
+ exitDescription, exitDescription
|
||||
String exitDescription = stepExecution.getExitStatus().getExitDescription();
|
||||
assertTrue("Exit description does not contain exception type name: " + exitDescription, exitDescription
|
||||
.contains(AsynchronousFailureException.class.getName()));
|
||||
|
||||
}
|
||||
@@ -300,24 +281,22 @@ public class ChunkMessageItemWriterIntegrationTests {
|
||||
/**
|
||||
* @param expected
|
||||
* @param maxWait
|
||||
* @throws InterruptedException
|
||||
* @throws InterruptedException
|
||||
*/
|
||||
private void waitForResults(int expected, int maxWait) throws InterruptedException {
|
||||
int count = 0;
|
||||
while (TestItemWriter.count<expected && count<maxWait) {
|
||||
while (TestItemWriter.count < expected && count < maxWait) {
|
||||
count++;
|
||||
Thread.sleep(10);
|
||||
}
|
||||
}
|
||||
|
||||
private StepExecution getStepExecution(Step step)
|
||||
throws JobExecutionAlreadyRunningException, JobRestartException,
|
||||
private StepExecution getStepExecution(Step step) throws JobExecutionAlreadyRunningException, JobRestartException,
|
||||
JobInstanceAlreadyCompleteException {
|
||||
SimpleJob job = new SimpleJob();
|
||||
job.setName("job");
|
||||
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(),
|
||||
new JobParametersBuilder().addLong("job.counter", jobCounter++)
|
||||
.toJobParameters());
|
||||
JobExecution jobExecution = jobRepository.createJobExecution(job.getName(), new JobParametersBuilder().addLong(
|
||||
"job.counter", jobCounter++).toJobParameters());
|
||||
StepExecution stepExecution = jobExecution.createStepExecution(step.getName());
|
||||
return stepExecution;
|
||||
}
|
||||
|
||||
@@ -18,7 +18,6 @@ package org.springframework.batch.integration.job;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
import java.lang.annotation.Annotation;
|
||||
import java.lang.reflect.Method;
|
||||
@@ -31,6 +30,7 @@ import org.springframework.batch.core.JobInstance;
|
||||
import org.springframework.batch.core.JobParameters;
|
||||
import org.springframework.batch.core.StepExecution;
|
||||
import org.springframework.batch.integration.JobRepositorySupport;
|
||||
import org.springframework.batch.repeat.ExitStatus;
|
||||
import org.springframework.beans.factory.annotation.Required;
|
||||
import org.springframework.core.annotation.AnnotationUtils;
|
||||
import org.springframework.integration.channel.DirectChannel;
|
||||
@@ -101,17 +101,14 @@ public class MessageOrientedStepTests {
|
||||
public void onMessage(Message<?> message) {
|
||||
}
|
||||
});
|
||||
try {
|
||||
step.setExecutionTimeout(1000);
|
||||
step.setPollingInterval(100);
|
||||
step.execute(jobExecution.createStepExecution(step.getName()));
|
||||
fail("Expected StepExecutionTimeoutException");
|
||||
}
|
||||
catch (StepExecutionTimeoutException e) {
|
||||
// expected
|
||||
String message = e.getMessage();
|
||||
assertTrue("Wrong message: " + message, message.contains("waiting for steps"));
|
||||
}
|
||||
step.setExecutionTimeout(1000);
|
||||
step.setPollingInterval(100);
|
||||
StepExecution stepExecution = jobExecution.createStepExecution(step.getName());
|
||||
step.execute(stepExecution);
|
||||
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
|
||||
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
|
||||
String message = stepExecution.getExitStatus().getExitDescription();
|
||||
assertTrue("Wrong message: " + message, message.contains("StepExecutionTimeoutException"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -135,15 +132,12 @@ public class MessageOrientedStepTests {
|
||||
replyChannel.send(message);
|
||||
}
|
||||
});
|
||||
try {
|
||||
step.execute(jobExecution.createStepExecution(step.getName()));
|
||||
fail("Expected RuntimeException");
|
||||
}
|
||||
catch (RuntimeException e) {
|
||||
// expected
|
||||
String message = e.getMessage();
|
||||
assertEquals("Wrong message: " + message, "Planned failure", message);
|
||||
}
|
||||
StepExecution stepExecution = jobExecution.createStepExecution(step.getName());
|
||||
step.execute(stepExecution);
|
||||
assertEquals(BatchStatus.FAILED, stepExecution.getStatus());
|
||||
assertEquals(ExitStatus.FAILED.getExitCode(), stepExecution.getExitStatus().getExitCode());
|
||||
String message = stepExecution.getExitStatus().getExitDescription();
|
||||
assertTrue("Wrong message: " + message, message.contains("Planned failure"));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -32,7 +32,7 @@ import org.springframework.integration.endpoint.SourcePoller;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.message.Message;
|
||||
import org.springframework.integration.message.MessageConsumer;
|
||||
import org.springframework.integration.message.PollableSource;
|
||||
import org.springframework.integration.message.MessageSource;
|
||||
import org.springframework.integration.scheduling.IntervalTrigger;
|
||||
import org.springframework.integration.scheduling.SimpleTaskScheduler;
|
||||
import org.springframework.integration.scheduling.TaskScheduler;
|
||||
@@ -80,7 +80,7 @@ public class PollableSourceRetryTests {
|
||||
processed.add((String) payload);
|
||||
}
|
||||
};
|
||||
PollableSource<Object> source = getPollableSource(list);
|
||||
MessageSource<Object> source = getPollableSource(list);
|
||||
MessageChannel target = getChannel(handler);
|
||||
SourcePoller trigger = getSourcePoller(source, target, transactionManager, 1);
|
||||
TaskScheduler scheduler = getSchedulerWithErrorHandler(trigger);
|
||||
@@ -109,7 +109,7 @@ public class PollableSourceRetryTests {
|
||||
throw new RuntimeException("Planned failure: " + payload);
|
||||
}
|
||||
};
|
||||
PollableSource<Object> source = getPollableSource(list);
|
||||
MessageSource<Object> source = getPollableSource(list);
|
||||
MessageChannel target = getChannel(handler);
|
||||
SourcePoller trigger = getSourcePoller(source, target, null, 1);
|
||||
TaskScheduler scheduler = getSchedulerWithErrorHandler(trigger);
|
||||
@@ -140,7 +140,7 @@ public class PollableSourceRetryTests {
|
||||
throw new RuntimeException("Planned failure: " + payload);
|
||||
}
|
||||
};
|
||||
PollableSource<Object> source = getPollableSource(list);
|
||||
MessageSource<Object> source = getPollableSource(list);
|
||||
MessageChannel target = getChannel(handler);
|
||||
SourcePoller trigger = getSourcePoller(source, target, transactionManager, 1);
|
||||
TaskScheduler scheduler = getSchedulerWithErrorHandler(trigger);
|
||||
@@ -176,7 +176,7 @@ public class PollableSourceRetryTests {
|
||||
}
|
||||
};
|
||||
|
||||
PollableSource<Object> source = getPollableSource(list);
|
||||
MessageSource<Object> source = getPollableSource(list);
|
||||
MessageChannel target = getChannel(handler);
|
||||
SourcePoller trigger = getSourcePoller(source, target, transactionManager, 1);
|
||||
TaskScheduler scheduler = getSchedulerWithErrorHandler(trigger);
|
||||
@@ -215,7 +215,7 @@ public class PollableSourceRetryTests {
|
||||
}
|
||||
};
|
||||
|
||||
PollableSource<Object> source = getPollableSource(list);
|
||||
MessageSource<Object> source = getPollableSource(list);
|
||||
MessageChannel target = getChannel(handler);
|
||||
SourcePoller trigger = getSourcePoller(source, target, null, 1);
|
||||
SourcePoller task = (SourcePoller) getProxy(trigger, SourcePoller.class, new Advice[] {
|
||||
@@ -262,7 +262,7 @@ public class PollableSourceRetryTests {
|
||||
}
|
||||
};
|
||||
|
||||
PollableSource<Object> source = getPollableSource(list);
|
||||
MessageSource<Object> source = getPollableSource(list);
|
||||
MessageChannel target = getChannel(handler);
|
||||
// this was the old dispatch advice chain
|
||||
target = (MessageChannel) getProxy(target, MessageChannel.class,
|
||||
@@ -305,7 +305,7 @@ public class PollableSourceRetryTests {
|
||||
}
|
||||
};
|
||||
|
||||
PollableSource<Object> source = getPollableSource(list);
|
||||
MessageSource<Object> source = getPollableSource(list);
|
||||
MessageChannel target = getChannel(handler);
|
||||
|
||||
// this was the old dispatch advice chain
|
||||
@@ -334,7 +334,7 @@ public class PollableSourceRetryTests {
|
||||
|
||||
}
|
||||
|
||||
private SourcePoller getSourcePoller(PollableSource<Object> source, MessageChannel channel,
|
||||
private SourcePoller getSourcePoller(MessageSource<Object> source, MessageChannel channel,
|
||||
PlatformTransactionManager transactionManager, int maxMessagesPerPoll) {
|
||||
SourcePoller poller = new SourcePoller(source, channel, new IntervalTrigger(100));
|
||||
poller.setTransactionManager(transactionManager);
|
||||
@@ -358,7 +358,7 @@ public class PollableSourceRetryTests {
|
||||
lifecycle.stop();
|
||||
}
|
||||
|
||||
private PollableSource<Object> getPollableSource(List<String> list) {
|
||||
private MessageSource<Object> getPollableSource(List<String> list) {
|
||||
final ItemReader<String> reader = new ListItemReader<String>(list) {
|
||||
public String read() {
|
||||
String item = super.read();
|
||||
@@ -366,7 +366,7 @@ public class PollableSourceRetryTests {
|
||||
return item;
|
||||
}
|
||||
};
|
||||
PollableSource<Object> source = new PollableSource<Object>() {
|
||||
MessageSource<Object> source = new MessageSource<Object>() {
|
||||
public Message<Object> receive() {
|
||||
try {
|
||||
String payload = reader.read();
|
||||
|
||||
Reference in New Issue
Block a user