BATCHADM-49: Added redelivered flag to ChunkResponse

This commit is contained in:
David Syer
2010-04-19 16:36:12 +00:00
committed by Michael Minella
parent c87d9418f5
commit 933bb75ee3
8 changed files with 158 additions and 27 deletions

View File

@@ -1,9 +1,9 @@
#Fri Apr 09 18:01:26 BST 2010
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/ChunkStepIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1102" endstart\="1096" start\="901" startend\="967"/>\n<bounds height\="118" width\="79" x\="15" y\="17"/>\n</element>\n<element type\="job">\n<structure end\="1473" endstart\="1467" start\="1107" startend\="1182"/>\n<bounds height\="118" width\="127" x\="106" y\="17"/>\n</element>\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1667" endstart\="1661" start\="1306" startend\="1372"/>\n<bounds height\="118" width\="79" x\="15" y\="17"/>\n</element>\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJmsIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1608" endstart\="1602" start\="1257" startend\="1323"/>\n<bounds height\="118" width\="79" x\="15" y\="17"/>\n</element>\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkStepIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1269" endstart\="1263" start\="1067" startend\="1133"/>\n<bounds height\="118" width\="79" x\="15" y\="17"/>\n</element>\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/file/FileToMessagesJobIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\r\n<graph>\r\n<element type\="job">\r\n<structure end\="972" endstart\="966" start\="771" startend\="837"/>\r\n<bounds height\="118" width\="75" x\="15" y\="17"/>\r\n</element>\r\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/step/StepGatewayIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1419" endstart\="1413" start\="1302" startend\="1368"/>\n<bounds height\="128" width\="79" x\="109" y\="19"/>\n</element>\n<element type\="step">\n<structure end\="1953" endstart\="1946" start\="1846" startend\="1916"/>\n<bounds height\="34" width\="74" x\="17" y\="19"/>\n</element>\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/tasklet/StepGatewayIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1421" endstart\="1415" start\="1302" startend\="1368"/>\n<bounds height\="128" width\="79" x\="109" y\="19"/>\n</element>\n<element type\="step">\n<structure end\="1962" endstart\="1955" start\="1855" startend\="1925"/>\n<bounds height\="34" width\="74" x\="17" y\="19"/>\n</element>\n</graph>
eclipse.preferences.version=1
#Mon Apr 19 14:39:29 BST 2010
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/ChunkStepIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1102" endstart\="1096" start\="901" startend\="967"/>\n<bounds height\="118" width\="79" x\="15" y\="17"/>\n</element>\n<element type\="job">\n<structure end\="1473" endstart\="1467" start\="1107" startend\="1182"/>\n<bounds height\="118" width\="127" x\="106" y\="17"/>\n</element>\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1667" endstart\="1661" start\="1306" startend\="1372"/>\n<bounds height\="118" width\="79" x\="15" y\="17"/>\n</element>\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkFaultTolerantStepJmsIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\r\n<graph>\r\n<element type\="job">\r\n<structure end\="1767" endstart\="1761" start\="1416" startend\="1482"/>\r\n<bounds height\="118" width\="75" x\="15" y\="17"/>\r\n</element>\r\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/chunk/RemoteChunkStepIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1269" endstart\="1263" start\="1067" startend\="1133"/>\n<bounds height\="118" width\="79" x\="15" y\="17"/>\n</element>\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/file/FileToMessagesJobIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\r\n<graph>\r\n<element type\="job">\r\n<structure end\="972" endstart\="966" start\="771" startend\="837"/>\r\n<bounds height\="118" width\="75" x\="15" y\="17"/>\r\n</element>\r\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/step/StepGatewayIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1419" endstart\="1413" start\="1302" startend\="1368"/>\n<bounds height\="128" width\="79" x\="109" y\="19"/>\n</element>\n<element type\="step">\n<structure end\="1953" endstart\="1946" start\="1846" startend\="1916"/>\n<bounds height\="34" width\="74" x\="17" y\="19"/>\n</element>\n</graph>
//com.springsource.sts.config.flow.coordinates\:http\://www.springframework.org/schema/batch\:/spring-batch-integration/src/test/resources/org/springframework/batch/integration/tasklet/StepGatewayIntegrationTests-context.xml=<?xml version\="1.0" encoding\="UTF-8"?>\n<graph>\n<element type\="job">\n<structure end\="1421" endstart\="1415" start\="1302" startend\="1368"/>\n<bounds height\="128" width\="79" x\="109" y\="19"/>\n</element>\n<element type\="step">\n<structure end\="1962" endstart\="1955" start\="1855" startend\="1925"/>\n<bounds height\="34" width\="74" x\="17" y\="19"/>\n</element>\n</graph>
eclipse.preferences.version=1

View File

@@ -1,8 +1,7 @@
<?xml version="1.0" encoding="UTF-8"?>
<project-modules id="moduleCoreId" project-version="1.5.0">
<wb-module deploy-name="spring-batch-integration">
<wb-resource deploy-path="/" source-path="/src/main/java"/>
<wb-resource deploy-path="/" source-path="/src/main/resources"/>
<wb-resource deploy-path="/" source-path="/src/test/resources"/>
</wb-module>
</project-modules>
<?xml version="1.0" encoding="UTF-8"?>
<project-modules id="moduleCoreId" project-version="1.5.0">
<wb-module deploy-name="spring-batch-integration">
<wb-resource deploy-path="/" source-path="/src/main/java"/>
<wb-resource deploy-path="/" source-path="/src/main/resources"/>
</wb-module>
</project-modules>

View File

@@ -95,7 +95,7 @@ public class ChunkMessageChannelItemWriter<T> extends StepExecutionListenerSuppo
ChunkRequest<T> request = new ChunkRequest<T>(items, localState.getJobId(), localState
.createStepContribution());
messagingGateway.send(request);
localState.expected.incrementAndGet();
localState.incrementExpected();
}
@@ -201,8 +201,17 @@ public class ChunkMessageChannelItemWriter<T> extends StepExecutionListenerSuppo
Assert.state(jobInstanceId != null, "Message did not contain job instance id.");
Assert.state(jobInstanceId.equals(localState.getJobId()), "Message contained wrong job instance id ["
+ jobInstanceId + "] should have been [" + localState.getJobId() + "].");
localState.actual.incrementAndGet();
if (payload.isRedelivered()) {
logger
.warn("Redelivered result detected, which may indicate stale state. In the best case, we just picked up a timed out message "
+ "from a previous failed execution. In the worst case (and if this is not a restart), "
+ "the step may now timeout. In that case if you believe that all messages "
+ "from workers have been sent, the business state "
+ "is probably inconsistent, and the step will fail.");
localState.incrementRedelivered();
}
localState.pushStepContribution(payload.getStepContribution());
localState.incrementActual();
if (!payload.isSuccessful()) {
throw new AsynchronousFailureException("Failure or interrupt detected in handler: "
+ payload.getMessage());
@@ -232,6 +241,8 @@ public class ChunkMessageChannelItemWriter<T> extends StepExecutionListenerSuppo
private AtomicInteger expected = new AtomicInteger();
private AtomicInteger redelivered = new AtomicInteger();
private StepExecution stepExecution;
private Queue<StepContribution> contributions = new LinkedBlockingQueue<StepContribution>();
@@ -263,6 +274,18 @@ public class ChunkMessageChannelItemWriter<T> extends StepExecutionListenerSuppo
}
}
public void incrementRedelivered() {
redelivered.incrementAndGet();
}
public void incrementActual() {
actual.incrementAndGet();
}
public void incrementExpected() {
expected.incrementAndGet();
}
public StepContribution createStepContribution() {
return stepExecution.createStepContribution();
}

View File

@@ -35,6 +35,8 @@ public class ChunkResponse implements Serializable {
private final boolean status;
private final String message;
private final boolean redelivered;
public ChunkResponse(Long jobId, StepContribution stepContribution) {
this(true, jobId, stepContribution, null);
@@ -45,10 +47,19 @@ public class ChunkResponse implements Serializable {
}
public ChunkResponse(boolean status, Long jobId, StepContribution stepContribution, String message) {
this(status, jobId, stepContribution, message, false);
}
public ChunkResponse(ChunkResponse input, boolean redelivered) {
this(input.status, input.jobId, input.stepContribution, input.message, redelivered);
}
public ChunkResponse(boolean status, Long jobId, StepContribution stepContribution, String message, boolean redelivered) {
this.status = status;
this.jobId = jobId;
this.stepContribution = stepContribution;
this.message = message;
this.redelivered = redelivered;
}
public StepContribution getStepContribution() {
@@ -62,6 +73,10 @@ public class ChunkResponse implements Serializable {
public boolean isSuccessful() {
return status;
}
public boolean isRedelivered() {
return redelivered;
}
public String getMessage() {
return message;

View File

@@ -0,0 +1,35 @@
/*
* Copyright 2009-2010 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.batch.integration.chunk;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.integration.annotation.Header;
/**
* @author Dave Syer
*
*/
public class JmsRedeliveredExtractor {
private static final Log logger = LogFactory.getLog(JmsRedeliveredExtractor.class);
public ChunkResponse extract(ChunkResponse input, @Header("springintegration_jms_redelivered") boolean redelivered) {
logger.debug("Extracted redelivered flag for response, value="+redelivered);
return new ChunkResponse(input, redelivered);
}
}

View File

@@ -26,6 +26,8 @@ public class MessageSourcePollerInterceptor extends ChannelInterceptorAdapter im
private MessageSource<?> source;
private MessageChannel channel;
/**
* Convenient default constructor for configuration purposes.
*/
@@ -39,6 +41,16 @@ public class MessageSourcePollerInterceptor extends ChannelInterceptorAdapter im
this.source = source;
}
/**
* Optional MessageChannel for injecting the message receieved from the source (defaults to the channel
* intercepted in {@link #preReceive(MessageChannel)}).
*
* @param channel the channel to set
*/
public void setChannel(MessageChannel channel) {
this.channel = channel;
}
/**
* Asserts that mandatory properties are set.
* @see InitializingBean#afterPropertiesSet()
@@ -64,6 +76,9 @@ public class MessageSourcePollerInterceptor extends ChannelInterceptorAdapter im
public boolean preReceive(MessageChannel channel) {
Message<?> message = source.receive();
if (message != null) {
if (this.channel!=null) {
channel = this.channel;
}
channel.send(message);
if (logger.isDebugEnabled()) {
logger.debug("Sent " + message + " to channel " + channel.getName());

View File

@@ -17,15 +17,52 @@ Remote Chunking Implementation
[[1]] Worker thread replies to response channel
[[1]] Step picks up reply and aggregates the counts
[[1]] Step picks up reply if there is one and aggregates the counts
[[1]] Step blocks until all the requests are satisfied
[[1]] Step reapeats until no more input data
[[1]] Step blocks until all the outstanding requests are satisfied
* Implementation
A ChunkProcessor acts as an asynchronous Messaging Gateway.
Asynchronous Messaging Gateway isn't a pattern that is supported out
of the box with Spring Integration. There is a
SimpleMessagingGateway that provides programmatic access to send and
receive payloads (instead of messages), so the pattern can be
manually implemented.
A ChunkProcessor acts as a kind of Throttling Asynchronous Messaging
Gateway, which isn't a pattern that is supported out of the box with
Spring Integration. There is a SimpleMessagingGateway that provides
programmatic access to send and receive payloads (instead of
messages), so the pattern can be manually implemented in the
ChunkProcessor.
The current implementation is in the form of an ItemWriter
(ChunkMessageChannelItemWriter) which is a StepExecutionListener
(blocks and waits for the outstanding responses in the afterStep).
The ChunkProcessor can then simply be a vanilla implementation from
Spring Batch.
The ChunkMessageChannelItemWriter implements the Throttling part of
the pattern by keeping track of the number of outstanding requests
(which it has to do anyway) and blocking until a response arrives if
the number is above a configurable limit. It wouldn't be necessary
to do this manually in the writer if the messages were only going
over local MessageChannels: the requests would either be processed
serially in a single thread, or else there would be a thread pool
with limited size controlling the workers. But since the messages
are going to JMS we need to either explicitly throttle in the writer
(or else rely on vendor features for producer flow control),
otherwise the JMS Queue could easily be overwhelmed and start
barfing (which happened in one of the early prototypes on an
Accenture project).
Throttling MessageChannel.send() might be something Spring
Integration could do, but it only makes sense really in the context
of this gateway pattern (because you need something to react against
to decide when to release another send).
The gateway is used to send requests to the workers, and then to
receive responses in the same thread, but only waiting for a
response when the step is complete. To implement this with JMS
backed channels we need a Spring Integration inbound adapter that
translates PollableChannel.receive() into
JmsTemplate.receiveAndConvert(). In the unlikely event of a problem
in the receiver the JMS message should roll back.
JmsDestinationPollingAdapter actually almost does what we need but
there is no support for configuring it without a scheduled poller.

View File

@@ -56,6 +56,12 @@
<int-jms:outbound-channel-adapter connection-factory="connectionFactory" channel="requests"
destination-name="requests" />
<integration:channel id="requests" />
<integration:channel id="incoming" />
<integration:transformer input-channel="incoming" output-channel="replies" ref="headerExtractor" method="extract"/>
<bean id="headerExtractor" class="org.springframework.batch.integration.chunk.JmsRedeliveredExtractor"/>
<integration:thread-local-channel id="replies">
<integration:interceptors>
<bean id="pollerInterceptor" class="org.springframework.batch.integration.chunk.MessageSourcePollerInterceptor">
@@ -70,6 +76,7 @@
</constructor-arg>
</bean>
</property>
<property name="channel" ref="incoming"/>
</bean>
</integration:interceptors>
</integration:thread-local-channel>