RESOLVED BATCH-1476. Count filters in chunk processor properly.
This commit is contained in:
@@ -137,7 +137,7 @@ public class FaultTolerantChunkProcessor<I, O> extends SimpleChunkProcessor<I, O
|
||||
@SuppressWarnings("unchecked")
|
||||
UserData<O> data = (UserData<O>) inputs.getUserData();
|
||||
if (data == null) {
|
||||
data = new UserData<O>(inputs.size());
|
||||
data = new UserData<O>();
|
||||
inputs.setUserData(data);
|
||||
data.setOutputs(new Chunk<O>());
|
||||
}
|
||||
@@ -147,7 +147,7 @@ public class FaultTolerantChunkProcessor<I, O> extends SimpleChunkProcessor<I, O
|
||||
protected int getFilterCount(Chunk<I> inputs, Chunk<O> outputs) {
|
||||
@SuppressWarnings("unchecked")
|
||||
UserData<O> data = (UserData<O>) inputs.getUserData();
|
||||
return data.size() - outputs.size() - inputs.getSkips().size();
|
||||
return data.filterCount;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -192,7 +192,7 @@ public class FaultTolerantChunkProcessor<I, O> extends SimpleChunkProcessor<I, O
|
||||
|
||||
Chunk<O> outputs = new Chunk<O>();
|
||||
@SuppressWarnings("unchecked")
|
||||
UserData<O> data = (UserData<O>) inputs.getUserData();
|
||||
final UserData<O> data = (UserData<O>) inputs.getUserData();
|
||||
final Chunk<O> cache = data.getOutputs();
|
||||
final Iterator<O> cacheIterator = cache.isEmpty() ? null : new ArrayList<O>(cache.getItems()).iterator();
|
||||
final AtomicInteger count = new AtomicInteger(0);
|
||||
@@ -253,6 +253,7 @@ public class FaultTolerantChunkProcessor<I, O> extends SimpleChunkProcessor<I, O
|
||||
if (output == null) {
|
||||
// No need to re-process filtered items
|
||||
iterator.remove();
|
||||
data.incrementFilterCount();
|
||||
}
|
||||
return output;
|
||||
}
|
||||
@@ -513,16 +514,12 @@ public class FaultTolerantChunkProcessor<I, O> extends SimpleChunkProcessor<I, O
|
||||
|
||||
private static class UserData<O> {
|
||||
|
||||
private final int size;
|
||||
|
||||
private Chunk<O> outputs;
|
||||
|
||||
public UserData(int size) {
|
||||
this.size = size;
|
||||
}
|
||||
private int filterCount = 0;
|
||||
|
||||
public int size() {
|
||||
return size;
|
||||
public void incrementFilterCount() {
|
||||
filterCount++;
|
||||
}
|
||||
|
||||
public Chunk<O> getOutputs() {
|
||||
|
||||
@@ -151,6 +151,7 @@ public class FaultTolerantChunkProcessorTests {
|
||||
}
|
||||
assertEquals(1, contribution.getSkipCount());
|
||||
assertEquals(1, contribution.getWriteCount());
|
||||
assertEquals(0, contribution.getFilterCount());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -391,6 +391,11 @@ public class FaultTolerantStepFactoryBeanRollbackTests {
|
||||
assertEquals("[1, 2, 3, 5]", writer.getCommitted().toString());
|
||||
assertEquals("[1, 2, 3, 4, 1, 2, 3, 4, 5]", writer.getWritten().toString());
|
||||
assertEquals("[1, 2, 3, 4, 5, 1, 2, 3, 4, 5]", processor.getProcessed().toString());
|
||||
|
||||
assertEquals(1, stepExecution.getWriteSkipCount());
|
||||
assertEquals(5, stepExecution.getReadCount());
|
||||
assertEquals(4, stepExecution.getWriteCount());
|
||||
assertEquals(0, stepExecution.getFilterCount());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -440,6 +445,11 @@ public class FaultTolerantStepFactoryBeanRollbackTests {
|
||||
assertEquals("[1, 2, 1, 2, 3, 4, 5]", writer.getWritten().toString());
|
||||
assertEquals("[1, 3, 5]", processor.getCommitted().toString());
|
||||
assertEquals("[1, 2, 3, 4, 5, 1, 2, 3, 4, 5]", processor.getProcessed().toString());
|
||||
|
||||
assertEquals(2, stepExecution.getWriteSkipCount());
|
||||
assertEquals(5, stepExecution.getReadCount());
|
||||
assertEquals(3, stepExecution.getWriteCount());
|
||||
assertEquals(0, stepExecution.getFilterCount());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user