Simplified the sum aggregation processor hierarchy

This commit is contained in:
Lennart Hensler committed 2014-02-26 15:16:13 +01:00
1 parent 1ddacd3946
commit 0aef8d3bc7
3 files changed
+39 -84

No files matched your search

@@ -1,47 +0,0 @@
package com.sap.sse.datamining.impl.components.aggregators;
import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.Map.Entry;
import java.util.concurrent.Executor;
import com.sap.sse.datamining.components.Processor;
import com.sap.sse.datamining.impl.components.GroupedDataEntry;
import com.sap.sse.datamining.shared.GroupKey;
public abstract class AbstractParallelGroupedDataSumAggregationProcessor<InputType, AggregatedType>
extends AbstractParallelSumAggregationProcessor<GroupedDataEntry<InputType>, Map<GroupKey, AggregatedType>> {
public AbstractParallelGroupedDataSumAggregationProcessor(Executor executor,
Collection<Processor<Map<GroupKey, AggregatedType>>> resultReceivers) {
super(executor, resultReceivers);
}
@Override
protected Map<GroupKey, AggregatedType> aggregateResult() {
Map<GroupKey, AggregatedType> result = new HashMap<>();
for (Entry<GroupedDataEntry<InputType>, Integer> elementAmountEntry : getElementAmountMap().entrySet()) {
InputType element = elementAmountEntry.getKey().getDataEntry();
Integer times = elementAmountEntry.getValue();
AggregatedType multipliedElementValue = multiply(element, times);
GroupKey groupKey = elementAmountEntry.getKey().getKey();
AggregatedType groupResult = result.get(groupKey);
result.put(groupKey, addToGroupResult(groupResult, multipliedElementValue));
}
return result;
}
private AggregatedType addToGroupResult(AggregatedType groupResult, AggregatedType multipliedElementValue) {
if (groupResult == null) {
return multipliedElementValue;
}
return add(groupResult, multipliedElementValue);
}
protected abstract AggregatedType multiply(InputType element, Integer times);
protected abstract AggregatedType add(AggregatedType firstSummand, AggregatedType secondSummand);
}
@@ -1,33 +0,0 @@
package com.sap.sse.datamining.impl.components.aggregators;
import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.Executor;
import com.sap.sse.datamining.components.Processor;
public abstract class AbstractParallelSumAggregationProcessor<InputType, AggregatedType>
extends AbstractStoringParallelAggregationProcessor<InputType, AggregatedType> {
private Map<InputType, Integer> elementAmountMap;
public AbstractParallelSumAggregationProcessor(Executor executor, Collection<Processor<AggregatedType>> resultReceivers) {
super(executor, resultReceivers);
elementAmountMap = new HashMap<>();
}
@Override
protected void storeElement(InputType element) {
if (!elementAmountMap.containsKey(element)) {
elementAmountMap.put(element, 0);
}
Integer currentAmount = elementAmountMap.get(element);
elementAmountMap.put(element, currentAmount + 1);
}
protected Map<InputType, Integer> getElementAmountMap() {
return elementAmountMap;
}
}
@@ -1,28 +1,63 @@
package com.sap.sse.datamining.impl.components.aggregators;
import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.Map.Entry;
import java.util.concurrent.Executor;
import com.sap.sse.datamining.components.Processor;
import com.sap.sse.datamining.impl.components.GroupedDataEntry;
import com.sap.sse.datamining.shared.GroupKey;
public class ParallelGroupedDoubleDataSumAggregationProcessor
extends AbstractParallelGroupedDataSumAggregationProcessor<Double, Double> {
extends AbstractStoringParallelAggregationProcessor<GroupedDataEntry<Double>, Map<GroupKey, Double>> {
private Map<GroupedDataEntry<Double>, Integer> elementAmountMap;
public ParallelGroupedDoubleDataSumAggregationProcessor(Executor executor,
Collection<Processor<Map<GroupKey, Double>>> resultReceivers) {
super(executor, resultReceivers);
elementAmountMap = new HashMap<>();
}
@Override
protected Double multiply(Double element, Integer times) {
protected void storeElement(GroupedDataEntry<Double> element) {
if (!elementAmountMap.containsKey(element)) {
elementAmountMap.put(element, 0);
}
Integer currentAmount = elementAmountMap.get(element);
elementAmountMap.put(element, currentAmount + 1);
}
@Override
protected Map<GroupKey, Double> aggregateResult() {
Map<GroupKey, Double> result = new HashMap<>();
for (Entry<GroupedDataEntry<Double>, Integer> elementAmountEntry : elementAmountMap.entrySet()) {
Double element = elementAmountEntry.getKey().getDataEntry();
Integer times = elementAmountEntry.getValue();
Double multipliedElementValue = multiply(element, times);
GroupKey groupKey = elementAmountEntry.getKey().getKey();
Double groupResult = result.get(groupKey);
result.put(groupKey, addToGroupResult(groupResult, multipliedElementValue));
}
return result;
}
private Double multiply(Double element, Integer times) {
return element * times;
}
@Override
protected Double add(Double firstSummand, Double secondSummand) {
private Double add(Double firstSummand, Double secondSummand) {
return firstSummand + secondSummand;
}
private Double addToGroupResult(Double groupResult, Double multipliedElementValue) {
if (groupResult == null) {
return multipliedElementValue;
}
return add(groupResult, multipliedElementValue);
}
}