Incomplete - task 87: Refactor KeyedItemReader
Split off from ItemReader interface.
This commit is contained in:
@@ -35,11 +35,11 @@ import org.springframework.batch.execution.step.support.ThreadStepInterruptionPo
|
||||
import org.springframework.batch.io.Skippable;
|
||||
import org.springframework.batch.io.exception.InfrastructureException;
|
||||
import org.springframework.batch.item.ExecutionContext;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.ItemRecoverer;
|
||||
import org.springframework.batch.item.ItemStream;
|
||||
import org.springframework.batch.item.ItemWriter;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.exception.CommitFailedException;
|
||||
import org.springframework.batch.repeat.ExitStatus;
|
||||
import org.springframework.batch.repeat.RepeatCallback;
|
||||
@@ -99,6 +99,26 @@ public class ItemOrientedStep extends AbstractStep implements InitializingBean {
|
||||
|
||||
private ListenerMulticaster listener = new ListenerMulticaster();
|
||||
|
||||
private ItemKeyGenerator itemKeyGenerator;
|
||||
|
||||
private ItemKeyGenerator defaultKeyGenerator = new ItemKeyGenerator() {
|
||||
public Object getKey(Object item) {
|
||||
return item;
|
||||
}
|
||||
};
|
||||
|
||||
/**
|
||||
* Public setter for the {@link ItemKeyGenerator}. If it is not injected
|
||||
* but the reader or writer implement {@link ItemKeyGenerator}, one of
|
||||
* those will be used instead (preferring the reader to the writer if both
|
||||
* would be appropriate).
|
||||
*
|
||||
* @param itemKeyGenerator the {@link ItemKeyGenerator} to set
|
||||
*/
|
||||
public void setItemKeyGenerator(ItemKeyGenerator itemKeyGenerator) {
|
||||
this.itemKeyGenerator = itemKeyGenerator;
|
||||
}
|
||||
|
||||
/**
|
||||
* Register each of the objects as listeners. The {@link ItemOrientedStep}
|
||||
* accepts listeners of type {@link ItemStream} and {@link BatchListener}.
|
||||
@@ -196,10 +216,10 @@ public class ItemOrientedStep extends AbstractStep implements InitializingBean {
|
||||
ItemReaderRetryPolicy itemProviderRetryPolicy = new ItemReaderRetryPolicy(retryPolicy);
|
||||
template.setRetryPolicy(itemProviderRetryPolicy);
|
||||
|
||||
itemKeyGenerator = getKeyGenerator();
|
||||
|
||||
if (retryPolicy != null) {
|
||||
Assert.state(itemReader instanceof KeyedItemReader,
|
||||
"ItemReader must be instance of KeyedItemReader to use the retry policy");
|
||||
retryCallback = new ItemReaderRetryCallback((KeyedItemReader) itemReader, itemWriter);
|
||||
retryCallback = new ItemReaderRetryCallback(itemReader, itemKeyGenerator, itemWriter);
|
||||
}
|
||||
|
||||
if (this.chunkOperations instanceof RepeatTemplate && commitInterval > 0) {
|
||||
@@ -212,6 +232,23 @@ public class ItemOrientedStep extends AbstractStep implements InitializingBean {
|
||||
|
||||
}
|
||||
|
||||
/**
|
||||
* @return
|
||||
*/
|
||||
private ItemKeyGenerator getKeyGenerator() {
|
||||
if (itemKeyGenerator != null) {
|
||||
return itemKeyGenerator;
|
||||
}
|
||||
if (itemReader instanceof ItemKeyGenerator) {
|
||||
return (ItemKeyGenerator) itemReader;
|
||||
}
|
||||
if (itemWriter instanceof ItemKeyGenerator) {
|
||||
return (ItemKeyGenerator) itemWriter;
|
||||
}
|
||||
return defaultKeyGenerator;
|
||||
|
||||
}
|
||||
|
||||
/**
|
||||
* Process the step and update its context so that progress can be monitored
|
||||
* by the caller. The step is broken down into chunks, each one executing in
|
||||
|
||||
@@ -30,8 +30,9 @@ import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.springframework.batch.io.Skippable;
|
||||
import org.springframework.batch.item.ExecutionContext;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.ItemStream;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.exception.ResetFailedException;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.dao.DataAccessException;
|
||||
@@ -107,7 +108,7 @@ import org.springframework.util.StringUtils;
|
||||
* @author Lucas Ward
|
||||
* @author Peter Zozom
|
||||
*/
|
||||
public class JdbcCursorItemReader implements KeyedItemReader, InitializingBean,
|
||||
public class JdbcCursorItemReader implements ItemReader, ItemKeyGenerator, InitializingBean,
|
||||
ItemStream, Skippable {
|
||||
|
||||
private static Log log = LogFactory.getLog(JdbcCursorItemReader.class);
|
||||
|
||||
@@ -19,8 +19,9 @@ import java.util.Iterator;
|
||||
import java.util.List;
|
||||
|
||||
import org.springframework.batch.item.ExecutionContext;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.ItemStream;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
@@ -46,7 +47,7 @@ import org.springframework.util.Assert;
|
||||
* @author Lucas Ward
|
||||
* @since 1.0
|
||||
*/
|
||||
public class DrivingQueryItemReader implements KeyedItemReader, InitializingBean,
|
||||
public class DrivingQueryItemReader implements ItemReader, ItemKeyGenerator, InitializingBean,
|
||||
ItemStream {
|
||||
|
||||
private boolean initialized = false;
|
||||
|
||||
@@ -22,7 +22,7 @@ package org.springframework.batch.item;
|
||||
* @author Dave Syer
|
||||
*
|
||||
*/
|
||||
public interface KeyedItemReader extends ItemReader {
|
||||
public interface ItemKeyGenerator {
|
||||
|
||||
/**
|
||||
* Get a unique identifier for the item that can be used to cache it between
|
||||
@@ -16,14 +16,14 @@
|
||||
package org.springframework.batch.item.reader;
|
||||
|
||||
import org.springframework.batch.item.ItemRecoverer;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
|
||||
|
||||
/**
|
||||
* @author Dave Syer
|
||||
*
|
||||
*/
|
||||
public abstract class AbstractItemReaderRecoverer extends AbstractItemReader implements KeyedItemReader, ItemRecoverer {
|
||||
public abstract class AbstractItemReaderRecoverer extends AbstractItemReader implements ItemKeyGenerator, ItemRecoverer {
|
||||
public Object getKey(Object item) {
|
||||
return item;
|
||||
}
|
||||
|
||||
@@ -18,7 +18,7 @@ package org.springframework.batch.item.reader;
|
||||
|
||||
import org.springframework.batch.io.Skippable;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
@@ -29,7 +29,7 @@ import org.springframework.util.Assert;
|
||||
*
|
||||
* @author Dave Syer
|
||||
*/
|
||||
public class DelegatingItemReader extends AbstractItemReader implements Skippable, InitializingBean, KeyedItemReader {
|
||||
public class DelegatingItemReader extends AbstractItemReader implements Skippable, InitializingBean, ItemKeyGenerator {
|
||||
|
||||
private ItemReader itemReader;
|
||||
|
||||
|
||||
@@ -24,7 +24,7 @@ import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.springframework.batch.item.FailedItemIdentifier;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.batch.item.exception.UnexpectedInputException;
|
||||
import org.springframework.jms.JmsException;
|
||||
import org.springframework.jms.core.JmsOperations;
|
||||
@@ -40,7 +40,7 @@ import org.springframework.util.Assert;
|
||||
* @author Dave Syer
|
||||
*
|
||||
*/
|
||||
public class JmsItemReader extends AbstractItemReader implements KeyedItemReader, FailedItemIdentifier {
|
||||
public class JmsItemReader extends AbstractItemReader implements ItemKeyGenerator, FailedItemIdentifier {
|
||||
|
||||
protected Log logger = LogFactory.getLog(getClass());
|
||||
|
||||
|
||||
@@ -18,10 +18,10 @@ package org.springframework.batch.retry.callback;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.ItemRecoverer;
|
||||
import org.springframework.batch.item.ItemWriter;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.retry.RetryCallback;
|
||||
import org.springframework.batch.retry.RetryContext;
|
||||
import org.springframework.batch.retry.RetryPolicy;
|
||||
@@ -42,27 +42,40 @@ import org.springframework.batch.retry.policy.ItemReaderRetryPolicy;
|
||||
*/
|
||||
public class ItemReaderRetryCallback implements RetryCallback {
|
||||
|
||||
private final static Log logger = LogFactory
|
||||
.getLog(ItemReaderRetryCallback.class);
|
||||
private final static Log logger = LogFactory.getLog(ItemReaderRetryCallback.class);
|
||||
|
||||
public static final String ITEM = ItemReaderRetryCallback.class.getName()
|
||||
+ ".ITEM";
|
||||
public static final String ITEM = ItemReaderRetryCallback.class.getName() + ".ITEM";
|
||||
|
||||
private KeyedItemReader provider;
|
||||
private ItemReader reader;
|
||||
|
||||
private ItemWriter writer;
|
||||
|
||||
private ItemRecoverer recoverer;
|
||||
|
||||
public ItemReaderRetryCallback(KeyedItemReader provider,
|
||||
ItemWriter writer) {
|
||||
private ItemKeyGenerator keyGenerator;
|
||||
|
||||
private ItemKeyGenerator defaultKeyGenerator = new ItemKeyGenerator() {
|
||||
public Object getKey(Object item) {
|
||||
return item;
|
||||
}
|
||||
};
|
||||
|
||||
public ItemReaderRetryCallback(ItemReader reader, ItemWriter writer) {
|
||||
this(reader, null, writer);
|
||||
}
|
||||
|
||||
public ItemReaderRetryCallback(ItemReader reader, ItemKeyGenerator keyGenerator, ItemWriter writer) {
|
||||
super();
|
||||
this.provider = provider;
|
||||
this.reader = reader;
|
||||
this.writer = writer;
|
||||
this.keyGenerator = keyGenerator;
|
||||
}
|
||||
|
||||
/**
|
||||
* Setter for injecting optional recovery handler.
|
||||
* Setter for injecting optional recovery handler. If it is not injected but
|
||||
* the reader or writer implement {@link ItemRecoverer}, one of those will
|
||||
* be used instead (preferring the reader to the writer if both would be
|
||||
* appropriate).
|
||||
*
|
||||
* @param recoveryHandler
|
||||
*/
|
||||
@@ -70,6 +83,17 @@ public class ItemReaderRetryCallback implements RetryCallback {
|
||||
this.recoverer = recoverer;
|
||||
}
|
||||
|
||||
/**
|
||||
* Public setter for the {@link ItemKeyGenerator}. If it is not injected
|
||||
* but the reader or writer implement {@link ItemKeyGenerator}, one of
|
||||
* those will be used instead (preferring the reader to the writer if both
|
||||
* would be appropriate).
|
||||
* @param keyGenerator the keyGenerator to set
|
||||
*/
|
||||
public void setKeyGenerator(ItemKeyGenerator keyGenerator) {
|
||||
this.keyGenerator = keyGenerator;
|
||||
}
|
||||
|
||||
public Object doWithRetry(RetryContext context) throws Throwable {
|
||||
// This requires a collaboration with the RetryPolicy...
|
||||
if (!context.isExhaustedOnly()) {
|
||||
@@ -82,10 +106,10 @@ public class ItemReaderRetryCallback implements RetryCallback {
|
||||
Object item = context.getAttribute(ITEM);
|
||||
if (item == null) {
|
||||
try {
|
||||
item = provider.read();
|
||||
} catch (Exception e) {
|
||||
throw new ExhaustedRetryException(
|
||||
"Unexpected end of item provider", e);
|
||||
item = reader.read();
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new ExhaustedRetryException("Unexpected end of item provider", e);
|
||||
}
|
||||
if (item == null) {
|
||||
// This is probably not fatal: in a batch we want to
|
||||
@@ -107,9 +131,31 @@ public class ItemReaderRetryCallback implements RetryCallback {
|
||||
}
|
||||
|
||||
/**
|
||||
* Accessor for the {@link ItemRecoverer}. If the handler is null but
|
||||
* the {@link ItemReader} is an instanceof {@link ItemRecoverer},
|
||||
* then it will be returned instead.
|
||||
* Accessor for the {@link ItemRecoverer}. If the handler is null but the
|
||||
* {@link ItemReader} is an instance of {@link ItemRecoverer}, then it will
|
||||
* be returned instead. If none of those strategies works then a default
|
||||
* implementation of {@link ItemKeyGenerator} will be used that just returns
|
||||
* the item.
|
||||
*
|
||||
* @return the {@link ItemRecoverer}.
|
||||
*/
|
||||
public ItemKeyGenerator getKeyGenerator() {
|
||||
if (keyGenerator != null) {
|
||||
return keyGenerator;
|
||||
}
|
||||
if (reader instanceof ItemKeyGenerator) {
|
||||
return (ItemKeyGenerator) reader;
|
||||
}
|
||||
if (writer instanceof ItemKeyGenerator) {
|
||||
return (ItemKeyGenerator) writer;
|
||||
}
|
||||
return defaultKeyGenerator;
|
||||
}
|
||||
|
||||
/**
|
||||
* Accessor for the {@link ItemRecoverer}. If the handler is null but the
|
||||
* {@link ItemReader} is an instance of {@link ItemRecoverer}, then it will
|
||||
* be returned instead.
|
||||
*
|
||||
* @return the {@link ItemRecoverer}.
|
||||
*/
|
||||
@@ -117,8 +163,11 @@ public class ItemReaderRetryCallback implements RetryCallback {
|
||||
if (recoverer != null) {
|
||||
return recoverer;
|
||||
}
|
||||
if (provider instanceof ItemRecoverer) {
|
||||
return (ItemRecoverer) provider;
|
||||
if (reader instanceof ItemRecoverer) {
|
||||
return (ItemRecoverer) reader;
|
||||
}
|
||||
if (writer instanceof ItemRecoverer) {
|
||||
return (ItemRecoverer) writer;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
@@ -128,8 +177,8 @@ public class ItemReaderRetryCallback implements RetryCallback {
|
||||
*
|
||||
* @return the {@link ItemReader} instance.
|
||||
*/
|
||||
public KeyedItemReader getReader() {
|
||||
return provider;
|
||||
public ItemReader getReader() {
|
||||
return reader;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -19,8 +19,9 @@ package org.springframework.batch.retry.policy;
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.springframework.batch.item.FailedItemIdentifier;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.ItemRecoverer;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.repeat.synch.RepeatSynchronizationManager;
|
||||
import org.springframework.batch.retry.RetryCallback;
|
||||
import org.springframework.batch.retry.RetryContext;
|
||||
@@ -32,11 +33,11 @@ import org.springframework.batch.retry.synch.RetrySynchronizationManager;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
* A {@link RetryPolicy} that detects an {@link ItemReaderRetryCallback} when
|
||||
* it opens a new context, and uses it to make sure the item is in place for
|
||||
* later decisions about how to retry or backoff. The callback should be an
|
||||
* instance of {@link ItemReaderRetryCallback} otherwise an exception will be
|
||||
* thrown when the context is created.
|
||||
* A {@link RetryPolicy} that detects an {@link ItemReaderRetryCallback} when it
|
||||
* opens a new context, and uses it to make sure the item is in place for later
|
||||
* decisions about how to retry or backoff. The callback should be an instance
|
||||
* of {@link ItemReaderRetryCallback} otherwise an exception will be thrown when
|
||||
* the context is created.
|
||||
*
|
||||
* @author Dave Syer
|
||||
*
|
||||
@@ -45,9 +46,7 @@ public class ItemReaderRetryPolicy extends AbstractStatefulRetryPolicy {
|
||||
|
||||
protected Log logger = LogFactory.getLog(getClass());
|
||||
|
||||
public static final String EXHAUSTED = ItemReaderRetryPolicy.class
|
||||
.getName()
|
||||
+ ".EXHAUSTED";
|
||||
public static final String EXHAUSTED = ItemReaderRetryPolicy.class.getName() + ".EXHAUSTED";
|
||||
|
||||
private RetryPolicy delegate;
|
||||
|
||||
@@ -103,14 +102,12 @@ public class ItemReaderRetryPolicy extends AbstractStatefulRetryPolicy {
|
||||
*
|
||||
* @see org.springframework.batch.retry.RetryPolicy#open(org.springframework.batch.retry.RetryCallback)
|
||||
*
|
||||
* @throws IllegalStateException
|
||||
* if the callback is not of the required type.
|
||||
* @throws IllegalStateException if the callback is not of the required
|
||||
* type.
|
||||
*/
|
||||
public RetryContext open(RetryCallback callback) {
|
||||
Assert.state(callback instanceof ItemReaderRetryCallback,
|
||||
"Callback must be ItemProviderRetryCallback");
|
||||
ItemReaderRetryContext context = new ItemReaderRetryContext(
|
||||
(ItemReaderRetryCallback) callback);
|
||||
Assert.state(callback instanceof ItemReaderRetryCallback, "Callback must be ItemProviderRetryCallback");
|
||||
ItemReaderRetryContext context = new ItemReaderRetryContext((ItemReaderRetryCallback) callback);
|
||||
context.open(callback);
|
||||
return context;
|
||||
}
|
||||
@@ -120,10 +117,9 @@ public class ItemReaderRetryPolicy extends AbstractStatefulRetryPolicy {
|
||||
* implemented by subclasses), and remove the current item from the history.
|
||||
*
|
||||
* @see org.springframework.batch.retry.RetryPolicy#registerThrowable(org.springframework.batch.retry.RetryContext,
|
||||
* java.lang.Throwable)
|
||||
* java.lang.Throwable)
|
||||
*/
|
||||
public void registerThrowable(RetryContext context, Throwable throwable)
|
||||
throws TerminatedRetryException {
|
||||
public void registerThrowable(RetryContext context, Throwable throwable) throws TerminatedRetryException {
|
||||
((RetryPolicy) context).registerThrowable(context, throwable);
|
||||
// The throwable is stored in the delegate context.
|
||||
}
|
||||
@@ -137,23 +133,25 @@ public class ItemReaderRetryPolicy extends AbstractStatefulRetryPolicy {
|
||||
return ((RetryPolicy) context).handleRetryExhausted(context);
|
||||
}
|
||||
|
||||
private class ItemReaderRetryContext extends RetryContextSupport
|
||||
implements RetryPolicy {
|
||||
private class ItemReaderRetryContext extends RetryContextSupport implements RetryPolicy {
|
||||
|
||||
private Object item;
|
||||
|
||||
// The delegate context...
|
||||
private RetryContext delegateContext;
|
||||
|
||||
private KeyedItemReader reader;
|
||||
private ItemReader reader;
|
||||
|
||||
private ItemRecoverer recoverer;
|
||||
|
||||
private ItemKeyGenerator keyGenerator;
|
||||
|
||||
public ItemReaderRetryContext(ItemReaderRetryCallback callback) {
|
||||
super(RetrySynchronizationManager.getContext());
|
||||
item = callback.next(this);
|
||||
this.reader = callback.getReader();
|
||||
this.recoverer = callback.getRecoverer();
|
||||
this.keyGenerator = callback.getKeyGenerator();
|
||||
}
|
||||
|
||||
public boolean canRetry(RetryContext context) {
|
||||
@@ -165,9 +163,8 @@ public class ItemReaderRetryPolicy extends AbstractStatefulRetryPolicy {
|
||||
}
|
||||
|
||||
public RetryContext open(RetryCallback callback) {
|
||||
if (hasFailed(reader, item)) {
|
||||
this.delegateContext = retryContextCache.get(reader
|
||||
.getKey(item));
|
||||
if (hasFailed(reader, keyGenerator, item)) {
|
||||
this.delegateContext = retryContextCache.get(keyGenerator.getKey(item));
|
||||
}
|
||||
if (this.delegateContext == null) {
|
||||
// Only create a new context if we don't know the history of
|
||||
@@ -178,37 +175,31 @@ public class ItemReaderRetryPolicy extends AbstractStatefulRetryPolicy {
|
||||
return null;
|
||||
}
|
||||
|
||||
public void registerThrowable(RetryContext context, Throwable throwable)
|
||||
throws TerminatedRetryException {
|
||||
retryContextCache.put(reader.getKey(item), this.delegateContext);
|
||||
public void registerThrowable(RetryContext context, Throwable throwable) throws TerminatedRetryException {
|
||||
retryContextCache.put(keyGenerator.getKey(item), this.delegateContext);
|
||||
delegate.registerThrowable(this.delegateContext, throwable);
|
||||
}
|
||||
|
||||
public boolean isExternal() {
|
||||
// Not called...
|
||||
throw new UnsupportedOperationException(
|
||||
"Not supported - this code should be unreachable.");
|
||||
throw new UnsupportedOperationException("Not supported - this code should be unreachable.");
|
||||
}
|
||||
|
||||
public boolean shouldRethrow(RetryContext context) {
|
||||
// Not called...
|
||||
throw new UnsupportedOperationException(
|
||||
"Not supported - this code should be unreachable.");
|
||||
throw new UnsupportedOperationException("Not supported - this code should be unreachable.");
|
||||
}
|
||||
|
||||
public Object handleRetryExhausted(RetryContext context)
|
||||
throws Exception {
|
||||
public Object handleRetryExhausted(RetryContext context) throws Exception {
|
||||
// If there is no going back, then we can remove the history
|
||||
retryContextCache.remove(reader.getKey(item));
|
||||
retryContextCache.remove(keyGenerator.getKey(item));
|
||||
RepeatSynchronizationManager.setCompleteOnly();
|
||||
if (recoverer != null) {
|
||||
boolean success = recoverer.recover(item, context
|
||||
.getLastThrowable());
|
||||
boolean success = recoverer.recover(item, context.getLastThrowable());
|
||||
if (!success) {
|
||||
int count = context.getRetryCount();
|
||||
logger.error(
|
||||
"Could not recover from error after retry exhausted after ["
|
||||
+ count + "] attempts.", context.getLastThrowable());
|
||||
logger.error("Could not recover from error after retry exhausted after [" + count + "] attempts.",
|
||||
context.getLastThrowable());
|
||||
}
|
||||
}
|
||||
return item;
|
||||
@@ -235,15 +226,16 @@ public class ItemReaderRetryPolicy extends AbstractStatefulRetryPolicy {
|
||||
* decision is delegated to the provider. Otherwise we just check the cache
|
||||
* for the item key.
|
||||
*
|
||||
* @param provider
|
||||
* @param reader
|
||||
* @param keyGenerator
|
||||
* @param item
|
||||
* @return
|
||||
*/
|
||||
protected boolean hasFailed(KeyedItemReader provider, Object item) {
|
||||
if (provider instanceof FailedItemIdentifier) {
|
||||
return ((FailedItemIdentifier) provider).hasFailed(item);
|
||||
protected boolean hasFailed(ItemReader reader, ItemKeyGenerator keyGenerator, Object item) {
|
||||
if (reader instanceof FailedItemIdentifier) {
|
||||
return ((FailedItemIdentifier) reader).hasFailed(item);
|
||||
}
|
||||
return retryContextCache.containsKey(provider.getKey(item));
|
||||
return retryContextCache.containsKey(keyGenerator.getKey(item));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -19,10 +19,10 @@ package org.springframework.batch.retry;
|
||||
import java.util.List;
|
||||
|
||||
import org.springframework.batch.item.ItemRecoverer;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.batch.item.reader.ListItemReader;
|
||||
|
||||
public class ListItemReaderRecoverer extends ListItemReader implements KeyedItemReader, ItemRecoverer {
|
||||
public class ListItemReaderRecoverer extends ListItemReader implements ItemKeyGenerator, ItemRecoverer {
|
||||
|
||||
/**
|
||||
* Delegate to super class constructor.
|
||||
@@ -47,7 +47,7 @@ public class ListItemReaderRecoverer extends ListItemReader implements KeyedItem
|
||||
/**
|
||||
* Return the item (assume it is its own key).
|
||||
*
|
||||
* @see org.springframework.batch.item.KeyedItemReader#getKey(java.lang.Object)
|
||||
* @see org.springframework.batch.item.ItemKeyGenerator#getKey(java.lang.Object)
|
||||
*/
|
||||
public Object getKey(Object item) {
|
||||
return item;
|
||||
|
||||
@@ -156,7 +156,7 @@ public class ItemReaderRetryCallbackTests extends TestCase {
|
||||
}
|
||||
|
||||
public void testGetKey() throws Exception {
|
||||
assertEquals("key0", callback.getReader().getKey("foo"));
|
||||
assertEquals("key0", callback.getKeyGenerator().getKey("foo"));
|
||||
}
|
||||
|
||||
public void testRecoverWithoutSession() throws Exception {
|
||||
|
||||
@@ -24,7 +24,8 @@ import java.util.List;
|
||||
import junit.framework.TestCase;
|
||||
|
||||
import org.springframework.batch.item.FailedItemIdentifier;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.reader.ListItemReader;
|
||||
import org.springframework.batch.item.writer.AbstractItemWriter;
|
||||
import org.springframework.batch.repeat.RepeatContext;
|
||||
@@ -41,17 +42,23 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
|
||||
private ItemReaderRetryPolicy policy = new ItemReaderRetryPolicy();
|
||||
|
||||
private org.springframework.batch.item.KeyedItemReader provider;
|
||||
private ItemReader reader;
|
||||
|
||||
private int count = 0;
|
||||
|
||||
private List list = new ArrayList();
|
||||
|
||||
private ItemKeyGenerator keyGenerator = new ItemKeyGenerator() {
|
||||
public Object getKey(Object item) {
|
||||
return item;
|
||||
}
|
||||
};
|
||||
|
||||
protected void setUp() throws Exception {
|
||||
super.setUp();
|
||||
// The list simulates a failed delivery, redelivery of the same message,
|
||||
// then a new message...
|
||||
provider = new ListItemReaderRecoverer(Arrays.asList(new String[] { "foo", "foo", "bar" })) {
|
||||
reader = new ListItemReaderRecoverer(Arrays.asList(new String[] { "foo", "foo", "bar" })) {
|
||||
public boolean recover(Object data, Throwable cause) {
|
||||
count++;
|
||||
list.add(data);
|
||||
@@ -61,7 +68,7 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
}
|
||||
|
||||
public void testOpenSunnyDay() throws Exception {
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
count++;
|
||||
list.add(data);
|
||||
@@ -71,8 +78,8 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
// we haven't called the processor yet...
|
||||
assertEquals(0, count);
|
||||
// but the provider has been accessed:
|
||||
assertEquals("foo", provider.read());
|
||||
assertEquals("bar", provider.read());
|
||||
assertEquals("foo", reader.read());
|
||||
assertEquals("bar", reader.read());
|
||||
}
|
||||
|
||||
public void testOpenWithWrongCallbackType() {
|
||||
@@ -92,7 +99,7 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
public void testCanRetry() {
|
||||
policy.setDelegate(new AlwaysRetryPolicy());
|
||||
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
count++;
|
||||
}
|
||||
@@ -105,7 +112,7 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
|
||||
public void testRegisterThrowable() {
|
||||
policy.setDelegate(new NeverRetryPolicy());
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
count++;
|
||||
list.add(data);
|
||||
@@ -118,7 +125,7 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
|
||||
public void testClose() throws Exception {
|
||||
policy.setDelegate(new NeverRetryPolicy());
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
count++;
|
||||
list.add(data);
|
||||
@@ -132,12 +139,12 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
// (not that this would happen in practice)...
|
||||
assertFalse(policy.canRetry(context));
|
||||
// The provider has been accessed only once:
|
||||
assertEquals("foo", provider.read());
|
||||
assertEquals("bar", provider.read());
|
||||
assertEquals("foo", reader.read());
|
||||
assertEquals("bar", reader.read());
|
||||
}
|
||||
|
||||
public void testOpenTwice() throws Exception {
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
count++;
|
||||
list.add(data);
|
||||
@@ -162,13 +169,13 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
// The provider has been accessed twice, so this
|
||||
// mimics a message receive by repeating the value of the first
|
||||
// message...
|
||||
assertEquals("bar", provider.read());
|
||||
assertEquals("bar", reader.read());
|
||||
}
|
||||
|
||||
public void testRecover() throws Exception {
|
||||
policy = new ItemReaderRetryPolicy();
|
||||
policy.setDelegate(new SimpleRetryPolicy(1));
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
}
|
||||
});
|
||||
@@ -207,7 +214,7 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
public void testRecoverWithTemplate() throws Exception {
|
||||
policy = new ItemReaderRetryPolicy();
|
||||
policy.setDelegate(new SimpleRetryPolicy(1));
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
throw new RuntimeException("Barf!");
|
||||
}
|
||||
@@ -230,7 +237,7 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
}
|
||||
|
||||
public void testExhaustedClearsHistoryAfterLastAttempt() throws Exception {
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
count++;
|
||||
list.add(data);
|
||||
@@ -258,7 +265,7 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
public void testRetryCount() throws Exception {
|
||||
policy = new ItemReaderRetryPolicy();
|
||||
policy.setDelegate(new SimpleRetryPolicy(1));
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
RetryContext context = policy.open(new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
count++;
|
||||
list.add(data);
|
||||
@@ -273,7 +280,7 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
}
|
||||
|
||||
public void testRetryCountPreservedBetweenRetries() throws Exception {
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(provider, new AbstractItemWriter() {
|
||||
ItemReaderRetryCallback callback = new ItemReaderRetryCallback(reader, new AbstractItemWriter() {
|
||||
public void write(Object data) {
|
||||
count++;
|
||||
list.add(data);
|
||||
@@ -296,10 +303,10 @@ public class ItemReaderRetryPolicyTests extends TestCase {
|
||||
MapRetryContextCache cache = new MapRetryContextCache();
|
||||
policy.setRetryContextCache(cache);
|
||||
cache.put("foo", new RetryContextSupport(null));
|
||||
assertTrue(policy.hasFailed(provider, "foo"));
|
||||
assertTrue(policy.hasFailed(reader, keyGenerator , "foo"));
|
||||
}
|
||||
|
||||
private static class MockFailedItemProvider extends ListItemReader implements KeyedItemReader, FailedItemIdentifier {
|
||||
private static class MockFailedItemProvider extends ListItemReader implements ItemKeyGenerator, FailedItemIdentifier {
|
||||
|
||||
private int hasFailedCount = 0;
|
||||
|
||||
|
||||
@@ -19,6 +19,7 @@ package org.springframework.batch.container.jms;
|
||||
import java.lang.reflect.InvocationTargetException;
|
||||
import java.lang.reflect.Method;
|
||||
|
||||
import javax.jms.ConnectionFactory;
|
||||
import javax.jms.JMSException;
|
||||
import javax.jms.Message;
|
||||
import javax.jms.MessageConsumer;
|
||||
@@ -48,16 +49,24 @@ public class BatchMessageListenerContainerTests extends TestCase {
|
||||
return ExitStatus.CONTINUABLE; // means we can continue to operate, but no message is received
|
||||
}
|
||||
};
|
||||
container = new BatchMessageListenerContainer(template);
|
||||
container = getContainer(template);
|
||||
boolean received = doExecute(null, null);
|
||||
assertEquals(1, count);
|
||||
assertFalse("Message received", received);
|
||||
}
|
||||
|
||||
private BatchMessageListenerContainer getContainer(RepeatTemplate template) {
|
||||
MockControl connectionFactoryControl = MockControl.createControl(ConnectionFactory.class);
|
||||
ConnectionFactory connectionFactory = (ConnectionFactory) connectionFactoryControl.getMock();
|
||||
BatchMessageListenerContainer container = new BatchMessageListenerContainer(template);
|
||||
container.setConnectionFactory(connectionFactory);
|
||||
return container;
|
||||
}
|
||||
|
||||
public void testReceiveAndExecuteWithCallback() throws Exception {
|
||||
RepeatTemplate template = new RepeatTemplate();
|
||||
template.setCompletionPolicy(new SimpleCompletionPolicy(2));
|
||||
container = new BatchMessageListenerContainer(template);
|
||||
container = getContainer(template);
|
||||
|
||||
MockControl sessionControl = MockControl.createNiceControl(Session.class);
|
||||
MockControl consumerControl = MockControl.createControl(MessageConsumer.class);
|
||||
@@ -87,7 +96,7 @@ public class BatchMessageListenerContainerTests extends TestCase {
|
||||
public void testReceiveAndExecuteWithCallbackReturningNull() throws Exception {
|
||||
RepeatTemplate template = new RepeatTemplate();
|
||||
template.setCompletionPolicy(new SimpleCompletionPolicy(2));
|
||||
container = new BatchMessageListenerContainer(template);
|
||||
container = getContainer(template);
|
||||
|
||||
MockControl sessionControl = MockControl.createNiceControl(Session.class);
|
||||
MockControl consumerControl = MockControl.createControl(MessageConsumer.class);
|
||||
@@ -114,7 +123,7 @@ public class BatchMessageListenerContainerTests extends TestCase {
|
||||
public void testTransactionalReceiveAndExecuteWithCallbackThrowingException() throws Exception {
|
||||
RepeatTemplate template = new RepeatTemplate();
|
||||
template.setCompletionPolicy(new SimpleCompletionPolicy(2));
|
||||
container = new BatchMessageListenerContainer(template);
|
||||
container = getContainer(template);
|
||||
container.setSessionTransacted(true);
|
||||
boolean received = doTestWithException(new IllegalStateException("No way!"), true, 2);
|
||||
assertFalse("Message received", received);
|
||||
@@ -123,7 +132,7 @@ public class BatchMessageListenerContainerTests extends TestCase {
|
||||
public void testNonTransactionalReceiveAndExecuteWithCallbackThrowingException() throws Exception {
|
||||
RepeatTemplate template = new RepeatTemplate();
|
||||
template.setCompletionPolicy(new SimpleCompletionPolicy(2));
|
||||
container = new BatchMessageListenerContainer(template);
|
||||
container = getContainer(template);
|
||||
container.setSessionTransacted(false);
|
||||
boolean received = doTestWithException(new IllegalStateException("No way!"), false, 2);
|
||||
assertTrue("Message not received", received);
|
||||
@@ -132,7 +141,7 @@ public class BatchMessageListenerContainerTests extends TestCase {
|
||||
public void testNonTransactionalReceiveAndExecuteWithCallbackThrowingError() throws Exception {
|
||||
RepeatTemplate template = new RepeatTemplate();
|
||||
template.setCompletionPolicy(new SimpleCompletionPolicy(2));
|
||||
container = new BatchMessageListenerContainer(template);
|
||||
container = getContainer(template);
|
||||
container.setSessionTransacted(false);
|
||||
try {
|
||||
boolean received = doTestWithException(new RuntimeException("No way!"), false, 2);
|
||||
|
||||
@@ -9,6 +9,7 @@ import junit.framework.TestCase;
|
||||
import org.springframework.batch.io.oxm.domain.Trade;
|
||||
import org.springframework.batch.io.xml.StaxEventItemReader;
|
||||
import org.springframework.batch.io.xml.oxm.UnmarshallingEventReaderDeserializer;
|
||||
import org.springframework.batch.item.ExecutionContext;
|
||||
import org.springframework.core.io.ClassPathResource;
|
||||
import org.springframework.core.io.Resource;
|
||||
import org.springframework.oxm.Unmarshaller;
|
||||
@@ -28,6 +29,9 @@ public abstract class AbstractStaxEventReaderItemReaderTests extends TestCase {
|
||||
source.setFragmentRootElementName("trade");
|
||||
UnmarshallingEventReaderDeserializer deserializer = new UnmarshallingEventReaderDeserializer(getUnmarshaller());
|
||||
source.setFragmentDeserializer(deserializer);
|
||||
|
||||
source.open(new ExecutionContext());
|
||||
|
||||
}
|
||||
|
||||
public void testRead() {
|
||||
|
||||
@@ -14,6 +14,7 @@ import org.custommonkey.xmlunit.XMLUnit;
|
||||
import org.springframework.batch.io.oxm.domain.Trade;
|
||||
import org.springframework.batch.io.xml.StaxEventItemWriter;
|
||||
import org.springframework.batch.io.xml.oxm.MarshallingEventWriterSerializer;
|
||||
import org.springframework.batch.item.ExecutionContext;
|
||||
import org.springframework.core.io.ClassPathResource;
|
||||
import org.springframework.core.io.FileSystemResource;
|
||||
import org.springframework.core.io.Resource;
|
||||
@@ -58,6 +59,8 @@ public abstract class AbstractStaxEventWriterItemWriterTests extends TestCase {
|
||||
|
||||
MarshallingEventWriterSerializer mapper = new MarshallingEventWriterSerializer(getMarshaller());
|
||||
writer.setSerializer(mapper);
|
||||
|
||||
writer.open(new ExecutionContext());
|
||||
}
|
||||
|
||||
protected void tearDown() throws Exception {
|
||||
|
||||
@@ -21,7 +21,7 @@ import java.util.List;
|
||||
|
||||
import javax.sql.DataSource;
|
||||
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.reader.AbstractItemReaderRecoverer;
|
||||
import org.springframework.batch.item.writer.AbstractItemWriter;
|
||||
import org.springframework.batch.repeat.ExitStatus;
|
||||
@@ -48,7 +48,7 @@ public class ExternalRetryInBatchTests extends AbstractDependencyInjectionSpring
|
||||
|
||||
private RepeatTemplate repeatTemplate;
|
||||
|
||||
private KeyedItemReader provider;
|
||||
private ItemReader provider;
|
||||
|
||||
private JdbcTemplate jdbcTemplate;
|
||||
|
||||
|
||||
@@ -21,7 +21,7 @@ import java.util.List;
|
||||
|
||||
import javax.sql.DataSource;
|
||||
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.reader.AbstractItemReaderRecoverer;
|
||||
import org.springframework.batch.item.writer.AbstractItemWriter;
|
||||
import org.springframework.batch.retry.callback.ItemReaderRetryCallback;
|
||||
@@ -41,7 +41,7 @@ public class ExternalRetryTests extends AbstractDependencyInjectionSpringContext
|
||||
|
||||
private RetryTemplate retryTemplate;
|
||||
|
||||
private KeyedItemReader provider;
|
||||
private ItemReader provider;
|
||||
|
||||
private JdbcTemplate jdbcTemplate;
|
||||
|
||||
|
||||
@@ -12,8 +12,9 @@ import org.apache.commons.logging.LogFactory;
|
||||
import org.springframework.batch.core.domain.StepExecution;
|
||||
import org.springframework.batch.core.domain.StepListener;
|
||||
import org.springframework.batch.item.ExecutionContext;
|
||||
import org.springframework.batch.item.ItemKeyGenerator;
|
||||
import org.springframework.batch.item.ItemReader;
|
||||
import org.springframework.batch.item.ItemStream;
|
||||
import org.springframework.batch.item.KeyedItemReader;
|
||||
import org.springframework.batch.item.exception.StreamException;
|
||||
import org.springframework.batch.repeat.ExitStatus;
|
||||
import org.springframework.batch.sample.item.writer.StagingItemWriter;
|
||||
@@ -25,7 +26,7 @@ import org.springframework.jdbc.support.lob.LobHandler;
|
||||
import org.springframework.transaction.support.TransactionSynchronizationManager;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
public class StagingItemReader extends JdbcDaoSupport implements ItemStream, KeyedItemReader, StepListener {
|
||||
public class StagingItemReader extends JdbcDaoSupport implements ItemStream, ItemReader, ItemKeyGenerator, StepListener {
|
||||
|
||||
// Key for buffer in transaction synchronization manager
|
||||
private static final String BUFFER_KEY = StagingItemReader.class.getName() + ".BUFFER";
|
||||
|
||||
Reference in New Issue
Block a user