generified the integration module and cleaned up the warnings
This commit is contained in:
@@ -1,8 +1,8 @@
|
||||
package org.springframework.batch.integration.chunk;
|
||||
|
||||
|
||||
public interface ChunkHandler {
|
||||
public interface ChunkHandler<T> {
|
||||
|
||||
ChunkResponse handleChunk(ChunkRequest chunk);
|
||||
ChunkResponse handleChunk(ChunkRequest<? extends T> chunk);
|
||||
|
||||
}
|
||||
@@ -22,7 +22,7 @@ import org.springframework.integration.message.Message;
|
||||
import org.springframework.transaction.support.TransactionSynchronizationManager;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
public class ChunkMessageChannelItemWriter extends StepExecutionListenerSupport implements ItemWriter, ItemStream {
|
||||
public class ChunkMessageChannelItemWriter<T> extends StepExecutionListenerSupport implements ItemWriter<T>, ItemStream {
|
||||
|
||||
private static final Log logger = LogFactory.getLog(ChunkMessageChannelItemWriter.class);
|
||||
|
||||
@@ -63,7 +63,7 @@ public class ChunkMessageChannelItemWriter extends StepExecutionListenerSupport
|
||||
this.requestChannel = requestChannel;
|
||||
}
|
||||
|
||||
public void write(Object item) throws Exception {
|
||||
public void write(T item) throws Exception {
|
||||
bindTransactionResources();
|
||||
getProcessed().add(item);
|
||||
logger.debug("Added item to chunk: " + item);
|
||||
@@ -88,13 +88,13 @@ public class ChunkMessageChannelItemWriter extends StepExecutionListenerSupport
|
||||
getNextResult(100);
|
||||
}
|
||||
|
||||
List<Object> processed = getProcessed();
|
||||
List<T> processed = getProcessed();
|
||||
|
||||
if (!processed.isEmpty()) {
|
||||
|
||||
logger.debug("Dispatching chunk: " + processed);
|
||||
ChunkRequest request = new ChunkRequest(processed, localState.getJobId(), localState.getSkipCount());
|
||||
GenericMessage<ChunkRequest> message = new GenericMessage<ChunkRequest>(request);
|
||||
ChunkRequest<T> request = new ChunkRequest<T>(processed, localState.getJobId(), localState.getSkipCount());
|
||||
GenericMessage<ChunkRequest<T>> message = new GenericMessage<ChunkRequest<T>>(request);
|
||||
requestChannel.send(message);
|
||||
localState.expected++;
|
||||
|
||||
@@ -197,11 +197,11 @@ public class ChunkMessageChannelItemWriter extends StepExecutionListenerSupport
|
||||
*
|
||||
* @return the processed
|
||||
*/
|
||||
@SuppressWarnings("unchecked")
|
||||
private List<Object> getProcessed() {
|
||||
private List<T> getProcessed() {
|
||||
Assert.state(TransactionSynchronizationManager.hasResource(ITEMS_PROCESSED),
|
||||
"Processed items not bound to transaction.");
|
||||
List<Object> processed = (List<Object>) TransactionSynchronizationManager.getResource(ITEMS_PROCESSED);
|
||||
@SuppressWarnings("unchecked")
|
||||
List<T> processed = (List<T>) TransactionSynchronizationManager.getResource(ITEMS_PROCESSED);
|
||||
return processed;
|
||||
}
|
||||
|
||||
|
||||
@@ -3,13 +3,13 @@ package org.springframework.batch.integration.chunk;
|
||||
import java.io.Serializable;
|
||||
import java.util.Collection;
|
||||
|
||||
public class ChunkRequest implements Serializable {
|
||||
public class ChunkRequest<T> implements Serializable {
|
||||
|
||||
private final int skipCount;
|
||||
private final Long jobId;
|
||||
private final Collection<Object> items;
|
||||
private final Collection<? extends T> items;
|
||||
|
||||
public ChunkRequest(Collection<Object> items, Long jobId, int skipCount) {
|
||||
public ChunkRequest(Collection<? extends T> items, Long jobId, int skipCount) {
|
||||
this.items = items;
|
||||
this.jobId = jobId;
|
||||
this.skipCount = skipCount;
|
||||
@@ -23,7 +23,7 @@ public class ChunkRequest implements Serializable {
|
||||
return jobId;
|
||||
}
|
||||
|
||||
public Collection<Object> getItems() {
|
||||
public Collection<? extends T> getItems() {
|
||||
return items;
|
||||
}
|
||||
|
||||
|
||||
@@ -11,11 +11,11 @@ import org.springframework.batch.repeat.ExitStatus;
|
||||
import org.springframework.integration.annotation.Handler;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
public class ItemWriterChunkHandler implements ChunkHandler {
|
||||
public class ItemWriterChunkHandler<T> implements ChunkHandler<T> {
|
||||
|
||||
private static final Log logger = LogFactory.getLog(ItemWriterChunkHandler.class);
|
||||
|
||||
private ItemWriter itemWriter;
|
||||
private ItemWriter<? super T> itemWriter;
|
||||
|
||||
private ItemSkipPolicy itemSkipPolicy = new NeverSkipItemSkipPolicy();
|
||||
|
||||
@@ -25,7 +25,7 @@ public class ItemWriterChunkHandler implements ChunkHandler {
|
||||
this.itemSkipPolicy = itemSkipPolicy;
|
||||
}
|
||||
|
||||
public void setItemWriter(ItemWriter itemWriter) {
|
||||
public void setItemWriter(ItemWriter<? super T> itemWriter) {
|
||||
this.itemWriter = itemWriter;
|
||||
}
|
||||
|
||||
@@ -45,7 +45,7 @@ public class ItemWriterChunkHandler implements ChunkHandler {
|
||||
*/
|
||||
@Handler
|
||||
@Transactional
|
||||
public ChunkResponse handleChunk(ChunkRequest chunk) {
|
||||
public ChunkResponse handleChunk(ChunkRequest<? extends T> chunk) {
|
||||
|
||||
logger.debug("Handling chunk: " + chunk);
|
||||
|
||||
@@ -53,7 +53,7 @@ public class ItemWriterChunkHandler implements ChunkHandler {
|
||||
int skipCount = 0;
|
||||
|
||||
try {
|
||||
for (Object item : chunk.getItems()) {
|
||||
for (T item : chunk.getItems()) {
|
||||
try {
|
||||
itemWriter.write(item);
|
||||
}
|
||||
|
||||
@@ -51,11 +51,11 @@ import org.springframework.util.Assert;
|
||||
* @author Dave Syer
|
||||
*
|
||||
*/
|
||||
public class FileToMessagesJobFactoryBean implements FactoryBean, BeanNameAware {
|
||||
public class FileToMessagesJobFactoryBean<T> implements FactoryBean, BeanNameAware {
|
||||
|
||||
private String name = "fileToMessageJob";
|
||||
|
||||
private ItemReader itemReader;
|
||||
private ItemReader<? extends T> itemReader;
|
||||
|
||||
private MessageChannel channel;
|
||||
|
||||
@@ -80,7 +80,7 @@ public class FileToMessagesJobFactoryBean implements FactoryBean, BeanNameAware
|
||||
* @param itemReader the itemReader to set
|
||||
*/
|
||||
@Required
|
||||
public void setItemReader(ItemReader itemReader) {
|
||||
public void setItemReader(ItemReader<? extends T> itemReader) {
|
||||
this.itemReader = itemReader;
|
||||
}
|
||||
|
||||
@@ -125,7 +125,7 @@ public class FileToMessagesJobFactoryBean implements FactoryBean, BeanNameAware
|
||||
job.setName(name);
|
||||
job.setJobRepository(jobRepository);
|
||||
|
||||
SimpleStepFactoryBean stepFactory = new SimpleStepFactoryBean();
|
||||
SimpleStepFactoryBean<T> stepFactory = new SimpleStepFactoryBean<T>();
|
||||
stepFactory.setBeanName("step");
|
||||
|
||||
Assert.state((itemReader instanceof FlatFileItemReader) || (itemReader instanceof StaxEventItemReader),
|
||||
@@ -139,7 +139,7 @@ public class FileToMessagesJobFactoryBean implements FactoryBean, BeanNameAware
|
||||
Assert.notNull(channel, "A channel must be provided");
|
||||
Assert.state(channel instanceof DirectChannel,
|
||||
"The channel must be a DirectChannel (otherwise failures can not be recovered from)");
|
||||
MessageChannelItemWriter itemWriter = new MessageChannelItemWriter();
|
||||
MessageChannelItemWriter<? super T> itemWriter = new MessageChannelItemWriter<T>();
|
||||
itemWriter.setChannel(channel);
|
||||
stepFactory.setItemWriter(itemWriter);
|
||||
|
||||
@@ -157,12 +157,12 @@ public class FileToMessagesJobFactoryBean implements FactoryBean, BeanNameAware
|
||||
* @param itemReader
|
||||
* @param resource
|
||||
*/
|
||||
private void setResource(ItemReader itemReader, Resource resource) {
|
||||
private void setResource(ItemReader<? extends T> itemReader, Resource resource) {
|
||||
if (itemReader instanceof FlatFileItemReader) {
|
||||
((FlatFileItemReader) itemReader).setResource(resource);
|
||||
((FlatFileItemReader<? extends T>) itemReader).setResource(resource);
|
||||
}
|
||||
else {
|
||||
((StaxEventItemReader) itemReader).setResource(resource);
|
||||
((StaxEventItemReader<? extends T>) itemReader).setResource(resource);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -24,7 +24,7 @@ import org.springframework.integration.message.GenericMessage;
|
||||
* @author Dave Syer
|
||||
*
|
||||
*/
|
||||
public class MessageChannelItemWriter extends AbstractItemWriter {
|
||||
public class MessageChannelItemWriter<T> extends AbstractItemWriter<T> {
|
||||
|
||||
private MessageChannel channel;
|
||||
|
||||
@@ -41,8 +41,8 @@ public class MessageChannelItemWriter extends AbstractItemWriter {
|
||||
* (non-Javadoc)
|
||||
* @see org.springframework.batch.item.ItemWriter#write(java.lang.Object)
|
||||
*/
|
||||
public void write(Object item) throws Exception {
|
||||
channel.send(new GenericMessage<Object>(item));
|
||||
public void write(T item) throws Exception {
|
||||
channel.send(new GenericMessage<T>(item));
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user