Implemented an abstract aggregation processor

This commit is contained in:
Lennart Hensler committed 2014-02-25 18:45:45 +01:00
1 parent 853ba51a02
commit f39cf9a32d
3 files changed
+118 -5

No files matched your search

@@ -0,0 +1,75 @@
package com.sap.sse.datamining.impl.components;
import static org.hamcrest.CoreMatchers.is;
import static org.junit.Assert.assertThat;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashSet;
import org.junit.Before;
import org.junit.Test;
import com.sap.sse.datamining.components.Processor;
import com.sap.sse.datamining.test.util.ConcurrencyTestsUtil;
public class TestAbstractStoringParallelAggregationProcessor {
private Collection<Processor<Integer>> receivers;
private boolean receiverWasToldToFinish = false;
private Integer receivedElement = null;
private Collection<Integer> elementStore = new ArrayList<>();
@Before
public void initializeReceivers() {
Processor<Integer> receiver = new Processor<Integer>() {
@Override
public void onElement(Integer element) {
receivedElement = element;
}
@Override
public void finish() throws InterruptedException {
receiverWasToldToFinish = true;
}
};
receivers = new HashSet<>();
receivers.add(receiver);
}
@Test
public void testAbstractAggregationHandling() throws InterruptedException {
Processor<Integer> processor = new AbstractStoringParallelAggregationProcessor<Integer, Integer>(ConcurrencyTestsUtil.getExecutor(), receivers) {
@Override
protected void storeElement(Integer element) {
elementStore.add(element);
}
@Override
protected Integer aggregateResult() {
Integer sum = 0;
for (Integer element : elementStore) {
sum += element;
}
return sum;
}
};
processElementAndVerifyThatItWasStored(processor, 42);
processElementAndVerifyThatItWasStored(processor, 7);
processor.finish();
ConcurrencyTestsUtil.sleepFor(100); //Giving the processor time to finish
assertThat("The receiver wasn't told to finish", receiverWasToldToFinish, is(true));
Integer expectedReceivedElement = 42 + 7;
assertThat(receivedElement, is(expectedReceivedElement));
}
private void processElementAndVerifyThatItWasStored(Processor<Integer> processor, int element) {
processor.onElement(element);
ConcurrencyTestsUtil.sleepFor(100); //Giving the processor time to process the instructions
assertThat("The element store doesn't contain the previously processed element '" + element + "'", elementStore.contains(element), is(true));
}
}
@@ -32,14 +32,15 @@ public abstract class AbstractPartitioningParallelProcessor<InputType, WorkingTy
@Override
public void onElement(InputType element) {
for (WorkingType partialElement : partitionElement(element)) {
final RunnableFuture<ResultType> instruction = new FutureTask<>(createInstruction(partialElement));
Callable<ResultType> instruction = createInstruction(partialElement);
if (isInstructionValid(instruction)) {
final RunnableFuture<ResultType> runnableInstruction = new FutureTask<>(instruction);
Runnable instructionWrapper = new Runnable() {
@Override
public void run() {
instruction.run();
runnableInstruction.run();
try {
ResultType result = instruction.get();
ResultType result = runnableInstruction.get();
if (isResultValid(result)) {
forwardResultToReceivers(result);
}
@@ -56,7 +57,7 @@ public abstract class AbstractPartitioningParallelProcessor<InputType, WorkingTy
}
}
protected boolean isInstructionValid(Runnable instruction) {
protected boolean isInstructionValid(Callable<ResultType> instruction) {
return instruction != null;
}
@@ -71,7 +72,7 @@ public abstract class AbstractPartitioningParallelProcessor<InputType, WorkingTy
return null;
}
private void forwardResultToReceivers(ResultType result) {
protected void forwardResultToReceivers(ResultType result) {
for (Processor<ResultType> resultReceiver : resultReceivers) {
resultReceiver.onElement(result);
}
@@ -0,0 +1,37 @@
package com.sap.sse.datamining.impl.components;
import java.util.Collection;
import java.util.concurrent.Callable;
import java.util.concurrent.Executor;
import com.sap.sse.datamining.components.Processor;
public abstract class AbstractStoringParallelAggregationProcessor<InputType, AggregatedType> extends
AbstractSimpleParallelProcessor<InputType, AggregatedType> {
public AbstractStoringParallelAggregationProcessor(Executor executor, Collection<Processor<AggregatedType>> resultReceivers) {
super(executor, resultReceivers);
}
@Override
protected Callable<AggregatedType> createInstruction(final InputType element) {
return new Callable<AggregatedType>() {
@Override
public AggregatedType call() throws Exception {
storeElement(element);
return AbstractStoringParallelAggregationProcessor.super.createInvalidResult();
}
};
}
protected abstract void storeElement(InputType element);
@Override
public void finish() throws InterruptedException {
super.forwardResultToReceivers(aggregateResult());
super.finish();
}
protected abstract AggregatedType aggregateResult();
}