Implemented a test for ProcessorQuery with standard workflow

This commit is contained in:
Lennart Hensler committed 2014-02-26 19:50:50 +01:00
1 parent 0fe488357d
commit 69b81f8a25
5 files changed
+263 -61

No files matched your search

@@ -0,0 +1,130 @@
package com.sap.sse.datamining.impl.components;
import static org.hamcrest.Matchers.is;
import static org.junit.Assert.assertThat;
import static org.junit.Assert.fail;
import java.lang.reflect.Method;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Map;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ThreadPoolExecutor;
import org.junit.Test;
import com.sap.sse.datamining.Query;
import com.sap.sse.datamining.components.Processor;
import com.sap.sse.datamining.factories.FunctionFactory;
import com.sap.sse.datamining.functions.Function;
import com.sap.sse.datamining.impl.components.aggregators.ParallelGroupedDoubleDataSumAggregationProcessor;
import com.sap.sse.datamining.shared.GroupKey;
import com.sap.sse.datamining.shared.QueryResult;
import com.sap.sse.datamining.shared.Unit;
import com.sap.sse.datamining.shared.impl.GenericGroupKey;
import com.sap.sse.datamining.shared.impl.QueryResultImpl;
import com.sap.sse.datamining.test.components.util.Number;
import com.sap.sse.datamining.test.util.ConcurrencyTestsUtil;
import com.sap.sse.datamining.test.util.FunctionTestsUtil;
public class TestProcessorQuery {
@Test
public void testStandardWorkflow() throws InterruptedException, ExecutionException {
Collection<Number> dataSource = createDataSource();
Query<Double> queryWithStandardWorkflow = createQueryWithStandardWorkflow(dataSource);
QueryResult<Double> expectedResult = buildExpectedResult(dataSource);
verifyResultWith(queryWithStandardWorkflow.run(), expectedResult);
}
private Collection<Number> createDataSource() {
Collection<Number> dataSource = new ArrayList<>();
//Will be filtered
dataSource.add(new Number(1));
dataSource.add(new Number(7));
//Results in <2> = 5
dataSource.add(new Number(10));
dataSource.add(new Number(10));
dataSource.add(new Number(10));
dataSource.add(new Number(10));
dataSource.add(new Number(10));
//Results in <3> = 3
dataSource.add(new Number(100));
dataSource.add(new Number(100));
dataSource.add(new Number(100));
//Results in <4> = 10
dataSource.add(new Number(1000));
dataSource.add(new Number(1000));
dataSource.add(new Number(1000));
dataSource.add(new Number(1000));
dataSource.add(new Number(1000));
dataSource.add(new Number(1000));
dataSource.add(new Number(1000));
dataSource.add(new Number(1000));
dataSource.add(new Number(1000));
dataSource.add(new Number(1000));
return dataSource;
}
/**
* Creates a query, that takes a Collection of Numbers, groups them by
* their length, extracts the cross sum and aggregates these as sum.
*/
private Query<Double> createQueryWithStandardWorkflow(Collection<Number> dataSource) {
ThreadPoolExecutor executor = ConcurrencyTestsUtil.getExecutor();
ProcessorQuery<Double, Iterable<Number>> query = new ProcessorQuery<Double, Iterable<Number>>(dataSource);
Collection<Processor<Map<GroupKey, Double>>> aggregationResultReceivers = asCollection(query.getResultReceiver());
Processor<GroupedDataEntry<Double>> sumAggregator =
new ParallelGroupedDoubleDataSumAggregationProcessor(executor, aggregationResultReceivers);
Collection<Processor<GroupedDataEntry<Double>>> extractionResultReceivers = asCollection(sumAggregator);
Method getCrossSumMethod = FunctionTestsUtil.getMethodFromClass(Number.class, "getCrossSum");
Function<Double> getCrossSumFunction = FunctionFactory.createMethodWrappingFunction(getCrossSumMethod);
Processor<GroupedDataEntry<Number>> crossSumExtractor = new ParallelGroupedElementsValueExtractionProcessor<Number, Double>(
executor, extractionResultReceivers, getCrossSumFunction);
Collection<Processor<GroupedDataEntry<Number>>> groupingResultReceivers = asCollection(crossSumExtractor);
Collection<Function<?>> dimensions = new ArrayList<>();
Function<Integer> getLengthFunction = FunctionFactory.createMethodWrappingFunction(FunctionTestsUtil.getMethodFromClass(Number.class, "getLength"));
dimensions.add(getLengthFunction);
Processor<Iterable<Number>> lengthGrouper = new ParallelMultiDimensionalGroupingProcessor<>(executor, groupingResultReceivers, dimensions);
query.setFirstProcessor(lengthGrouper);
return query;
}
private <T> Collection<T> asCollection(T value) {
Collection<T> collection = new ArrayList<>();
collection.add(value);
return collection;
}
private QueryResult<Double> buildExpectedResult(Collection<Number> dataSource) {
QueryResultImpl<Double> result = new QueryResultImpl<>(dataSource.size(), 2, "Cross sum (Sum)", Unit.None, 0);
result.addResult(new GenericGroupKey<Integer>(2), 5.0);
result.addResult(new GenericGroupKey<Integer>(3), 3.0);
result.addResult(new GenericGroupKey<Integer>(4), 10.0);
return result;
}
private void verifyResultWith(QueryResult<Double> result, QueryResult<Double> expectedResult) {
assertThat("Result values aren't correct.", result.getResults(), is(expectedResult.getResults()));
assertThat("Retrieved data amount isn't correct.", result.getRetrievedDataAmount(), is(expectedResult.getRetrievedDataAmount()));
assertThat("Filtered data amount isn't correct.", result.getFilteredDataAmount(), is(expectedResult.getFilteredDataAmount()));
assertThat("Result signifier isn't correct.", result.getResultSignifier(), is(expectedResult.getResultSignifier()));
assertThat("Unit isn't correct.", result.getUnit(), is(expectedResult.getUnit()));
assertThat("Value decimals aren't correct.", result.getValueDecimals(), is(expectedResult.getValueDecimals()));
}
@Test
public void testQueryTimeouting() {
fail("Not yet implemented");
}
}
@@ -1,16 +0,0 @@
package com.sap.sse.datamining.impl.components;
import static org.junit.Assert.*;
import org.junit.Ignore;
import org.junit.Test;
public class TestQueryProcessor {
@Ignore
@Test
public void test() {
fail("Not yet implemented");
}
}
@@ -0,0 +1,41 @@
package com.sap.sse.datamining.impl.components;
import java.util.ArrayList;
import java.util.Collection;
import java.util.concurrent.Callable;
import java.util.concurrent.Executor;
import com.sap.sse.datamining.components.FilterCriteria;
import com.sap.sse.datamining.components.Processor;
public abstract class AbstractFilteringRetrievalProcessor<InputType, WorkingType, ResultType>
extends AbstractPartitioningParallelProcessor<InputType, WorkingType, ResultType> {
private final FilterCriteria<WorkingType> criteria;
public AbstractFilteringRetrievalProcessor(Executor executor, Collection<Processor<ResultType>> resultReceivers, FilterCriteria<WorkingType> criteria) {
super(executor, resultReceivers);
this.criteria = criteria;
}
@Override
protected Iterable<WorkingType> partitionElement(InputType element) {
Collection<WorkingType> filteredData = new ArrayList<>();
Iterable<WorkingType> retrievedData = retrieveData(element);
for (WorkingType retrievedDataEntry : retrievedData) {
if (criteria.matches(retrievedDataEntry)) {
filteredData.add(retrievedDataEntry);
}
}
return filteredData;
}
protected abstract Iterable<WorkingType> retrieveData(InputType element);
//Redefinition of the method to set the parameter name to element instead of partial element.
//This makes the implementation of sub classes more fluent.
@Override
protected abstract Callable<ResultType> createInstruction(WorkingType filteredPartialElement);
}
@@ -0,0 +1,92 @@
package com.sap.sse.datamining.impl.components;
import java.util.Map;
import java.util.logging.Level;
import java.util.logging.Logger;
import com.sap.sse.datamining.Query;
import com.sap.sse.datamining.components.Processor;
import com.sap.sse.datamining.shared.GroupKey;
import com.sap.sse.datamining.shared.QueryResult;
public class ProcessorQuery<AggregatedType, DataSourceType> implements Query<AggregatedType> {
private static final Logger LOGGER = Logger.getLogger(ProcessorQuery.class.getSimpleName());
private final DataSourceType dataSource;
private Processor<DataSourceType> firstProcessor;
private final ProcessResultReceiver resultReceiver;
private final Object monitorObject = new Object();
private boolean workIsDone;
public ProcessorQuery(DataSourceType dataSource) {
this.dataSource = dataSource;
resultReceiver = new ProcessResultReceiver();
}
public void setFirstProcessor(Processor<DataSourceType> firstProcessor) {
this.firstProcessor = firstProcessor;
}
@Override
public QueryResult<AggregatedType> run() {
try {
return processQuery();
} catch (InterruptedException e) {
LOGGER.log(Level.WARNING, "The query processing got interrupted.", e);
}
return null;
}
private QueryResult<AggregatedType> processQuery() throws InterruptedException {
firstProcessor.onElement(dataSource);
firstProcessor.finish();
waitTillWorkIsDone();
return resultReceiver.getResult();
}
private void waitTillWorkIsDone() throws InterruptedException {
synchronized (monitorObject) {
while (!workIsDone) {
monitorObject.wait();
}
workIsDone = false;
}
}
Processor<Map<GroupKey, AggregatedType>> getResultReceiver() {
return resultReceiver;
}
private class ProcessResultReceiver implements Processor<Map<GroupKey, AggregatedType>> {
private QueryResult<AggregatedType> result;
@Override
public void onElement(Map<GroupKey, AggregatedType> groupedAggregations) {
result = constructResult(groupedAggregations);
}
private QueryResult<AggregatedType> constructResult(Map<GroupKey, AggregatedType> groupedAggregations) {
// TODO Auto-generated method stub
return null;
}
@Override
public void finish() throws InterruptedException {
synchronized (monitorObject) {
workIsDone = true;
monitorObject.notify();
}
}
public QueryResult<AggregatedType> getResult() {
return result;
}
}
}
@@ -1,45 +0,0 @@
package com.sap.sse.datamining.impl.components;
import java.util.Map;
import java.util.concurrent.ExecutionException;
import com.sap.sse.datamining.Query;
import com.sap.sse.datamining.components.Processor;
import com.sap.sse.datamining.shared.GroupKey;
import com.sap.sse.datamining.shared.QueryResult;
public class QueryProcessor<DataSourceType, AggregatedType> implements Query<AggregatedType>, Processor<Map<GroupKey, AggregatedType>> {
private final DataSourceType dataSource;
private Processor<DataSourceType> firstProcessor;
private QueryResult<AggregatedType> result;
public QueryProcessor(DataSourceType dataSource) {
this.dataSource = dataSource;
}
public void setFirstProcessor(Processor<DataSourceType> firstProcessor) {
this.firstProcessor = firstProcessor;
}
@Override
public QueryResult<AggregatedType> run() throws InterruptedException, ExecutionException {
firstProcessor.onElement(dataSource);
wait();
return result;
}
@Override
public void onElement(Map<GroupKey, AggregatedType> element) {
// TODO Auto-generated method stub
}
@Override
public void finish() throws InterruptedException {
// TODO Auto-generated method stub
}
}