Implemented the first version of an integration test aborting a query with heavy load instructions

This commit is contained in:
Lennart Hensler
2018-07-29 19:52:49 +02:00
parent 3c31932f38
commit cbc3222ad2
19 changed files with 350 additions and 52 deletions
@@ -0,0 +1,273 @@
package com.sap.sse.datamining.impl;
import static org.hamcrest.Matchers.is;
import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashSet;
import java.util.Map;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.FutureTask;
import java.util.concurrent.RunnableFuture;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import org.junit.Before;
import org.junit.Test;
import com.sap.sse.datamining.Query;
import com.sap.sse.datamining.QueryState;
import com.sap.sse.datamining.components.AdditionalResultDataBuilder;
import com.sap.sse.datamining.components.Processor;
import com.sap.sse.datamining.components.ProcessorInstruction;
import com.sap.sse.datamining.components.ProcessorInstructionHandler;
import com.sap.sse.datamining.data.QueryResult;
import com.sap.sse.datamining.impl.components.AbstractParallelProcessor;
import com.sap.sse.datamining.impl.components.AbstractProcessorInstruction;
import com.sap.sse.datamining.impl.components.AbstractRetrievalProcessor;
import com.sap.sse.datamining.impl.components.GroupedDataEntry;
import com.sap.sse.datamining.impl.components.ProcessorInstructionPriority;
import com.sap.sse.datamining.impl.components.aggregators.ParallelGroupedDataCollectingAsSetProcessor;
import com.sap.sse.datamining.shared.GroupKey;
import com.sap.sse.datamining.shared.data.QueryResultState;
import com.sap.sse.datamining.shared.impl.GenericGroupKey;
import com.sap.sse.datamining.test.util.components.StatefulBlockingInstruction;
public class TestAbortingHeavyLoadQuery {
private static final int ExecutorPoolSize = Math.min(3, Runtime.getRuntime().availableProcessors() - 2);
private static final int DataSourceSize = 2000;
private static final String GroupKeyPrefix = "G";
private static final long HeavyLoadInstructionDuration = 500;
private static final long AbortQueryDelay = (long) (HeavyLoadInstructionDuration * 1.1);
private static final boolean RecordExecution = true;
private ConcurrentLinkedQueue<String> executionRecord;
private ExecutorService executor;
private Query<HashSet<Integer>> query;
private Collection<Processor<?, ?>> processors;
@Test
public void testAbortingHeavyLoadQuery() throws InterruptedException, ExecutionException {
RunnableFuture<QueryResult<HashSet<Integer>>> queryTask = new FutureTask<>(() -> {
logExecution("Starting query execution");
long start = System.currentTimeMillis();
QueryResult<HashSet<Integer>> result = query.run();
long duration = System.currentTimeMillis() - start;
logExecution("Finished query in " + duration + "ms");
return result;
});
Thread worker = new Thread(queryTask, "Worker");
worker.start();
do {
Thread.sleep(10);
} while (query.getState() != QueryState.RUNNING);
Thread.sleep(AbortQueryDelay);
logExecution("Aborting query");
query.abort();
QueryResult<HashSet<Integer>> result = queryTask.get();
assertThat(result.getState(), is(QueryResultState.ABORTED));
assertTrue("The result is not empty", result.isEmpty());
for (Processor<?, ?> processor : processors) {
assertTrue("Processor wasn't aborted", processor.isAborted());
}
// TODO Ensure that no instructions are scheduled after the query has been aborted
// TODO More detailed assessment of the leftover instructions
executor.shutdown();
boolean terminated = executor.awaitTermination((long) (HeavyLoadInstructionDuration * 1.5), TimeUnit.MILLISECONDS);
if (RecordExecution) {
for (String string : executionRecord) {
System.out.println(string);
}
}
assertTrue("The executor didn't terminate in the given time", terminated);
}
private void logExecution(String message) {
if (RecordExecution) {
executionRecord.add(Thread.currentThread().getName() + ": " + message);
}
}
@Before
@SuppressWarnings("unchecked")
public void initialize() {
executor = new DataMiningExecutorService(ExecutorPoolSize);
// executor = new ThreadPoolExecutor(ExecutorPoolSize, ExecutorPoolSize, 0, TimeUnit.MILLISECONDS, new PriorityBlockingQueue<>());
// executor = new ThreadPoolExecutor(ExecutorPoolSize, ExecutorPoolSize, 0, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<>());
if (RecordExecution) {
executionRecord = new ConcurrentLinkedQueue<>();
}
processors = new ArrayList<>();
Class<Iterable<String>> dataSourceType = (Class<Iterable<String>>)(Class<?>) Iterable.class;
Class<GroupedDataEntry<Integer>> groupedType = (Class<GroupedDataEntry<Integer>>)(Class<?>) GroupedDataEntry.class;
Class<HashSet<Integer>> resultType = (Class<HashSet<Integer>>)(Class<?>) HashSet.class;
Collection<String> dataSource = new ArrayList<>(DataSourceSize);
for (int i = 0; i < DataSourceSize; i++) {
dataSource.add(GroupKeyPrefix + i);
}
query = new ProcessorQuery<HashSet<Integer>, Iterable<String>>(dataSource, resultType) {
@Override
protected Processor<Iterable<String>, ?> createChainAndReturnFirstProcessor(Processor<Map<GroupKey, HashSet<Integer>>, Void> resultReceiver) {
Processor<GroupedDataEntry<Integer>, Map<GroupKey, HashSet<Integer>>> aggregator = new ParallelGroupedDataCollectingAsSetProcessor<Integer>(executor, Collections.singleton(resultReceiver)) {
@Override
protected void storeElement(GroupedDataEntry<Integer> element) {
logExecution("Storing " + element);
super.storeElement(element);
}
};
Processor<Element, GroupedDataEntry<Integer>> grouper = new AbstractParallelProcessor<Element, GroupedDataEntry<Integer>>(Element.class, groupedType, executor, Collections.singleton(aggregator)) {
@Override
protected ProcessorInstruction<GroupedDataEntry<Integer>> createInstruction(Element element) {
return new AbstractProcessorInstruction<GroupedDataEntry<Integer>>(this, ProcessorInstructionPriority.Grouping) {
@Override
protected GroupedDataEntry<Integer> computeResult() throws Exception {
logExecution("Grouping " + element);
return new GroupedDataEntry<>(new GenericGroupKey<>(element.getName()), element.getValue());
}
};
}
@Override
protected void setAdditionalData(AdditionalResultDataBuilder additionalDataBuilder) { }
};
Processor<Element, Element> heavyLoadProcessor = new AbstractParallelProcessor<Element, Element>(Element.class, Element.class, executor, Collections.singleton(grouper)) {
@Override
protected ProcessorInstruction<Element> createInstruction(Element element) {
ProcessorInstruction<Element> instruction = new RecordingBlockingInstruction(this,
ProcessorInstructionPriority.Extraction, HeavyLoadInstructionDuration, element,
TestAbortingHeavyLoadQuery.this::logExecution);
return instruction;
}
@Override
protected void setAdditionalData(AdditionalResultDataBuilder additionalDataBuilder) { }
};
Processor<String, Element> retriever1 = new AbstractRetrievalProcessor<String, Element>(String.class, Element.class, executor, Collections.singleton(heavyLoadProcessor), 1) {
@Override
protected Iterable<Element> retrieveData(String element) {
logExecution("Retrieving " + ExecutorPoolSize + " elements from " + element);
Collection<Element> data = new ArrayList<>();
for (int i = 0; i < ExecutorPoolSize; i++) {
data.add(new Element(element, i));
}
return data;
}
};
Processor<Iterable<String>, String> retriever0 = new AbstractRetrievalProcessor<Iterable<String>, String>(dataSourceType, String.class, executor, Collections.singleton(retriever1), 0) {
@Override
protected Iterable<String> retrieveData(Iterable<String> element) {
logExecution("Retrieving data from data source");
return element;
}
};
processors.add(retriever0);
processors.add(retriever1);
processors.add(heavyLoadProcessor);
processors.add(grouper);
processors.add(aggregator);
return retriever0;
}
};
}
private static class RecordingBlockingInstruction extends StatefulBlockingInstruction<Element> {
private final Consumer<String> recorder;
public RecordingBlockingInstruction(ProcessorInstructionHandler<Element> handler,
ProcessorInstructionPriority priority, long blockDuration, Element result, Consumer<String> recorder) {
super(handler, priority, blockDuration, result);
this.recorder = recorder;
}
@Override
public void run() {
recorder.accept("Executing heavy load instruction for " + result);
super.run();
}
@Override
protected void actionBeforeBlock() {
recorder.accept("Starting work for heavy load instruction for " + result);
}
@Override
protected void actionAfterBlock() {
recorder.accept("Finished heavy load instruction for " + result);
}
}
private static class Element {
private final String name;
private final int value;
public Element(String name, int value) {
this.name = name;
this.value = value;
}
public String getName() {
return name;
}
public int getValue() {
return value;
}
@Override
public String toString() {
return getName() + "-" + getValue();
}
@Override
public int hashCode() {
final int prime = 31;
int result = 1;
result = prime * result + ((name == null) ? 0 : name.hashCode());
result = prime * result + value;
return result;
}
@Override
public boolean equals(Object obj) {
if (this == obj)
return true;
if (obj == null)
return false;
if (getClass() != obj.getClass())
return false;
Element other = (Element) obj;
if (name == null) {
if (other.name != null)
return false;
} else if (!name.equals(other.name))
return false;
if (value != other.value)
return false;
return true;
}
}
}
@@ -63,7 +63,7 @@ public class TestProcessorQuery {
Collection<Processor<Double, ?>> resultReceivers = new ArrayList<>();
resultReceivers.add(new AbortResultReceiver(resultReceiver));
return new BlockingProcessor<Iterable<Number>, Double>((Class<Iterable<Number>>)(Class<?>) Iterable.class, Double.class,
ConcurrencyTestsUtil.getExecutor(), resultReceivers, 1000) {
ConcurrencyTestsUtil.getSharedExecutor(), resultReceivers, 1000) {
@Override
protected Double createResult(Iterable<Number> element) {
return 0.0;
@@ -101,7 +101,7 @@ public class TestProcessorQuery {
Collection<Processor<Double, ?>> resultReceivers = new ArrayList<>();
resultReceivers.add(new AbortResultReceiver(resultReceiver));
return new BlockingProcessor<Iterable<Number>, Double>((Class<Iterable<Number>>)(Class<?>) Iterable.class, Double.class,
ConcurrencyTestsUtil.getExecutor(), resultReceivers, 1000) {
ConcurrencyTestsUtil.getSharedExecutor(), resultReceivers, 1000) {
@Override
protected Double createResult(Iterable<Number> element) {
return 0.0;
@@ -168,7 +168,7 @@ public class TestProcessorQuery {
resultReceivers.add(resultReceiver);
return new AbstractParallelProcessor<Iterable<Number>, Map<GroupKey, Double>>((Class<Iterable<Number>>)(Class<?>) Iterable.class,
(Class<Map<GroupKey, Double>>)(Class<?>) Map.class,
ConcurrencyTestsUtil.getExecutor(),
ConcurrencyTestsUtil.getSharedExecutor(),
resultReceivers) {
@Override
protected ProcessorInstruction<Map<GroupKey, Double>> createInstruction(final Iterable<Number> element) {
@@ -242,7 +242,7 @@ public class TestProcessorQuery {
resultReceivers.add(resultReceiver);
return new AbstractParallelProcessor<Double, Map<GroupKey, Double>>(Double.class,
(Class<Map<GroupKey, Double>>)(Class<?>) Map.class,
ConcurrencyTestsUtil.getExecutor(),
ConcurrencyTestsUtil.getSharedExecutor(),
resultReceivers) {
@Override
protected ProcessorInstruction<Map<GroupKey, Double>> createInstruction(Double element) {
@@ -282,7 +282,7 @@ public class TestProcessorQuery {
resultReceivers.add(resultReceiver);
return new AbstractParallelProcessor<Double, Map<GroupKey, Double>>(Double.class,
(Class<Map<GroupKey, Double>>)(Class<?>) Map.class,
ConcurrencyTestsUtil.getExecutor(),
ConcurrencyTestsUtil.getSharedExecutor(),
resultReceivers) {
@Override
protected ProcessorInstruction<Map<GroupKey, Double>> createInstruction(Double element) {
@@ -23,7 +23,7 @@ public class TestAbstractParallelProcessorElementProcessing {
private ThreadPoolExecutor executor;
private Processor<Integer, Object> processor;
private List<StatefulBlockingInstruction> createdInstructions;
private List<StatefulBlockingInstruction<?>> createdInstructions;
@Before
public void initialize() {
@@ -33,7 +33,7 @@ public class TestAbstractParallelProcessorElementProcessing {
processor = new AbstractParallelProcessor<Integer, Object>(Integer.class, Object.class, executor, Collections.emptySet()) {
@Override
protected ProcessorInstruction<Object> createInstruction(Integer sleepTime) {
StatefulBlockingInstruction instruction = new StatefulBlockingInstruction(this, sleepTime);
StatefulBlockingInstruction<Object> instruction = new StatefulBlockingInstruction<>(this, sleepTime);
createdInstructions.add(instruction);
return instruction;
}
@@ -52,7 +52,7 @@ public class TestAbstractParallelProcessorElementProcessing {
assertThat("Unexpected amount of created instructions", createdInstructions.size(), is(elementCount));
executor.shutdown();
assertThat("Executor couldn't terminate", executor.awaitTermination(1, TimeUnit.SECONDS), is(true));
for (StatefulBlockingInstruction instruction : createdInstructions) {
for (StatefulBlockingInstruction<?> instruction : createdInstructions) {
assertThat("run wasn't called", instruction.runWasCalled(), is(true));
assertThat("computeResult wasn't called", instruction.computeResultWasCalled(), is(true));
assertThat("computeResult didn't finish", instruction.computeResultWasFinished(), is(true));
@@ -84,7 +84,7 @@ public class TestAbstractParallelProcessorElementProcessing {
processor.processElement(instructionDuration);
assertThat("Unexpected amount of created instructions", createdInstructions.size(), is(elementCount));
for (StatefulBlockingInstruction instruction : createdInstructions) {
for (StatefulBlockingInstruction<?> instruction : createdInstructions) {
assertThat("run wasn't called", instruction.runWasCalled(), is(true));
assertThat("computeResult wasn't called", instruction.computeResultWasCalled(), is(true));
assertThat("computeResult didn't finish", instruction.computeResultWasFinished(), is(true));
@@ -108,7 +108,7 @@ public class TestAbstractParallelProcessorElementProcessing {
executor.shutdown();
assertThat("Executor couldn't terminate", executor.awaitTermination(1, TimeUnit.SECONDS), is(true));
for (int i = 0; i < createdInstructions.size(); i++) {
StatefulBlockingInstruction instruction = createdInstructions.get(i);
StatefulBlockingInstruction<?> instruction = createdInstructions.get(i);
assertThat("run wasn't called", instruction.runWasCalled(), is(true));
// Last instructions was executed after the processor has been aborted. computeResult() should not be called
if (i == createdInstructions.size() - 1) {
@@ -68,7 +68,7 @@ public class TestAbstractParallelProcessorFinishing {
}
private AbstractParallelProcessor<Integer, Integer> createProcessor(Collection<Processor<Integer, ?>> receivers) {
return new AbstractParallelProcessor<Integer, Integer>(Integer.class, Integer.class, ConcurrencyTestsUtil.getExecutor(), receivers) {
return new AbstractParallelProcessor<Integer, Integer>(Integer.class, Integer.class, ConcurrencyTestsUtil.getSharedExecutor(), receivers) {
@Override
protected ProcessorInstruction<Integer> createInstruction(Integer partialElement) {
return new AbstractProcessorInstruction<Integer>(this) {
@@ -37,7 +37,7 @@ public class TestAbstractParallelProcessorWithManySimpleInstructions {
Collection<Processor<Integer, ?>> receivers = new ArrayList<>();
receivers.add(receiver);
processor = new AbstractParallelProcessor<Integer, Integer>(Integer.class, Integer.class, ConcurrencyTestsUtil.getExecutor(), receivers) {
processor = new AbstractParallelProcessor<Integer, Integer>(Integer.class, Integer.class, ConcurrencyTestsUtil.getSharedExecutor(), receivers) {
@Override
protected ProcessorInstruction<Integer> createInstruction(final Integer element) {
return new AbstractProcessorInstruction<Integer>(this) {
@@ -1,6 +1,6 @@
package com.sap.sse.datamining.impl.components;
import static com.sap.sse.datamining.test.util.ConcurrencyTestsUtil.getExecutor;
import static com.sap.sse.datamining.test.util.ConcurrencyTestsUtil.getSharedExecutor;
import static org.hamcrest.Matchers.is;
import static org.hamcrest.Matchers.not;
import static org.junit.Assert.assertThat;
@@ -90,12 +90,12 @@ public class TestDataRetrieverChainCreation {
assertThat(chainClone, not(dataRetrieverChainDefinition));
chainClone.startBuilding(ConcurrencyTestsUtil.getExecutor());
chainClone.startBuilding(ConcurrencyTestsUtil.getSharedExecutor());
}
@Test
public void testThatTheDataRetrieverChainBuilderHasToBeInitialized() {
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getExecutor());
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getSharedExecutor());
try {
chainBuilder.getCurrentRetrievedDataType();
@@ -126,7 +126,7 @@ public class TestDataRetrieverChainCreation {
@Test
public void testStepByStepDataRetrieverChainCreation() throws InterruptedException {
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getExecutor());
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getSharedExecutor());
chainBuilder.stepFurther(); //Initialization
assertThat(chainBuilder.getCurrentRetrievedDataType().equals(Test_Regatta.class), is(true));
@@ -206,7 +206,7 @@ public class TestDataRetrieverChainCreation {
@Test(expected=IllegalArgumentException.class)
public void testSettingAFilterWithWrongElementType() {
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getExecutor());
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getSharedExecutor());
chainBuilder.stepFurther(); //Initialization
chainBuilder.setFilter(raceFilter);
@@ -214,7 +214,7 @@ public class TestDataRetrieverChainCreation {
@Test(expected=IllegalArgumentException.class)
public void testSettingAResultReceiverWithWrongInputType() {
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getExecutor());
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getSharedExecutor());
chainBuilder.stepFurther(); //Initialization
chainBuilder.addResultReceiver(legReceiver);
@@ -253,12 +253,12 @@ public class TestDataRetrieverChainCreation {
raceRetrieverClass,
Test_HasRaceContext.class, "race");
dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getExecutor());
dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getSharedExecutor());
}
@Test(expected=IllegalStateException.class)
public void testSteppingToFar() {
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getExecutor());
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getSharedExecutor());
while (chainBuilder.canStepFurther()) {
chainBuilder.stepFurther();
}
@@ -272,7 +272,7 @@ public class TestDataRetrieverChainCreation {
chainWithSettings.endWith(TestLegOfCompetitorWithContextRetrievalProcessor.class, Test_RetrievalProcessorWithSettings.class, Test_HasLegOfCompetitorContext.class,
Test_RetrievalProcessorSettings.class, new Test_RetrievalProcessorSettings("Default Settings"), "legOfCompetitor");
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = chainWithSettings.startBuilding(getExecutor());
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = chainWithSettings.startBuilding(getSharedExecutor());
while (chainBuilder.canStepFurther()) {
chainBuilder.stepFurther();
}
@@ -311,7 +311,7 @@ public class TestDataRetrieverChainCreation {
chainWithSettings.endWith(TestLegOfCompetitorWithContextRetrievalProcessor.class, Test_RetrievalProcessorWithSettings.class, Test_HasLegOfCompetitorContext.class,
Test_RetrievalProcessorSettings.class, new Test_RetrievalProcessorSettings("Default Settings"), "legOfCompetitor");
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = chainWithSettings.startBuilding(getExecutor());
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = chainWithSettings.startBuilding(getSharedExecutor());
while (chainBuilder.canStepFurther()) {
chainBuilder.stepFurther();
}
@@ -320,7 +320,7 @@ public class TestDataRetrieverChainCreation {
@Test(expected=IllegalStateException.class)
public void testRetrieverWithNoSettingsButSettedSettingsCreated() {
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getExecutor());
DataRetrieverChainBuilder<Collection<Test_Regatta>> chainBuilder = dataRetrieverChainDefinition.startBuilding(ConcurrencyTestsUtil.getSharedExecutor());
chainBuilder.stepFurther(); // Initialization
chainBuilder.setSettings(new WrongSettings());
}
@@ -34,7 +34,7 @@ public class TestFilteringProcessors {
return element % 2 == 0;
}
};
Processor<Integer, Integer> filteringProcessor = new ParallelFilteringProcessor<Integer>(Integer.class, ConcurrencyTestsUtil.getExecutor(), receivers, elementIsEvenCriteria);
Processor<Integer, Integer> filteringProcessor = new ParallelFilteringProcessor<Integer>(Integer.class, ConcurrencyTestsUtil.getSharedExecutor(), receivers, elementIsEvenCriteria);
ConcurrencyTestsUtil.processElements(filteringProcessor, createElementsToProcess());
ConcurrencyTestsUtil.sleepFor(100); // Giving the processor time to process the instructions
@@ -63,7 +63,7 @@ public class TestParallelExtractionProcessor {
@Test
public void testValueExtraction() {
Processor<GroupedDataEntry<Number>, GroupedDataEntry<Integer>> processor = new ParallelGroupedElementsValueExtractionProcessor<Number, Integer>(ConcurrencyTestsUtil.getExecutor(), receivers, getCrossSumFunction);
Processor<GroupedDataEntry<Number>, GroupedDataEntry<Integer>> processor = new ParallelGroupedElementsValueExtractionProcessor<Number, Integer>(ConcurrencyTestsUtil.getSharedExecutor(), receivers, getCrossSumFunction);
Collection<GroupedDataEntry<Number>> elements = createElements();
ConcurrencyTestsUtil.processElements(processor, elements);
ConcurrencyTestsUtil.sleepFor(100); //Giving the processor time to finish the instructions
@@ -81,7 +81,7 @@ public class TestParallelExtractionProcessor {
@Test
public void testValueExtractionWithInvalidFunction() {
Processor<GroupedDataEntry<Number>, GroupedDataEntry<Integer>> processor = new ParallelGroupedElementsValueExtractionProcessor<Number, Integer>(ConcurrencyTestsUtil.getExecutor(), receivers, invalidFunction);
Processor<GroupedDataEntry<Number>, GroupedDataEntry<Integer>> processor = new ParallelGroupedElementsValueExtractionProcessor<Number, Integer>(ConcurrencyTestsUtil.getSharedExecutor(), receivers, invalidFunction);
ConcurrencyTestsUtil.processElements(processor, createElements());
ConcurrencyTestsUtil.sleepFor(100); //Giving the processor time to finish the instructions
assertThat("Values have been received, but the processor function is invalid.", receivedValues.isEmpty(), is(true));
@@ -57,7 +57,7 @@ public class TestParallelMultiDimensionalGroupingProcessor {
@Test(expected=IllegalArgumentException.class)
public void testConstructionWithNullDimensions() {
new ParallelMultiDimensionsValueNestingGroupingProcessor<>(Number.class, ConcurrencyTestsUtil.getExecutor(), receivers, null);
new ParallelMultiDimensionsValueNestingGroupingProcessor<>(Number.class, ConcurrencyTestsUtil.getSharedExecutor(), receivers, null);
}
@Test(expected=IllegalArgumentException.class)
@@ -72,7 +72,7 @@ public class TestRetrieverFilterProcessorChain {
resultReceivers.add(retrievalProcessor);
@SuppressWarnings("unchecked")
Processor<Iterable<Iterable<Integer>>, Iterable<Integer>> layeredRetrievalProcessor = new AbstractRetrievalProcessor<Iterable<Iterable<Integer>>, Iterable<Integer>>((Class<Iterable<Iterable<Integer>>>)(Class<?>) Iterable.class, (Class<Iterable<Integer>>)(Class<?>) Iterable.class,
ConcurrencyTestsUtil.getExecutor(), resultReceivers, 0) {
ConcurrencyTestsUtil.getSharedExecutor(), resultReceivers, 0) {
@Override
protected Iterable<Iterable<Integer>> retrieveData(Iterable<Iterable<Integer>> element) {
return element;
@@ -110,11 +110,11 @@ public class TestRetrieverFilterProcessorChain {
return element >= 0;
}
};
Processor<Integer, Integer> filtrationProcessor = new ParallelFilteringProcessor<>(Integer.class, ConcurrencyTestsUtil.getExecutor(), filtrationResultReceivers, elementGreaterZeroFilterCriteria);
Processor<Integer, Integer> filtrationProcessor = new ParallelFilteringProcessor<>(Integer.class, ConcurrencyTestsUtil.getSharedExecutor(), filtrationResultReceivers, elementGreaterZeroFilterCriteria);
Collection<Processor<Integer, ?>> retrievalResultReceivers = new ArrayList<>();
retrievalResultReceivers.add(filtrationProcessor);
retrievalProcessor = new AbstractRetrievalProcessor<Iterable<Integer>, Integer>((Class<Iterable<Integer>>)(Class<?>) Iterable.class, Integer.class, ConcurrencyTestsUtil.getExecutor(), retrievalResultReceivers, 1) {
retrievalProcessor = new AbstractRetrievalProcessor<Iterable<Integer>, Integer>((Class<Iterable<Integer>>)(Class<?>) Iterable.class, Integer.class, ConcurrencyTestsUtil.getSharedExecutor(), retrievalResultReceivers, 1) {
@Override
protected Iterable<Integer> retrieveData(Iterable<Integer> element) {
return element;
@@ -48,7 +48,7 @@ public class TestAbstractStoringParallelAggregationProcessor {
@Test
public void testAbstractAggregationHandling() throws InterruptedException {
Processor<GroupedDataEntry<Integer>, Map<GroupKey, Integer>> processor = new AbstractParallelGroupedDataStoringAggregationProcessor<Integer, Integer>(ConcurrencyTestsUtil.getExecutor(), receivers, "Sum") {
Processor<GroupedDataEntry<Integer>, Map<GroupKey, Integer>> processor = new AbstractParallelGroupedDataStoringAggregationProcessor<Integer, Integer>(ConcurrencyTestsUtil.getSharedExecutor(), receivers, "Sum") {
@Override
protected void storeElement(GroupedDataEntry<Integer> element) {
elementStore.add(element);
@@ -104,7 +104,7 @@ public class TestAbstractStoringParallelAggregationProcessor {
}
}
});
Processor<GroupedDataEntry<Integer>, Map<GroupKey, Integer>> processor = new AbstractParallelGroupedDataStoringAggregationProcessor<Integer, Integer>(ConcurrencyTestsUtil.getExecutor(), receivers, "Sum") {
Processor<GroupedDataEntry<Integer>, Map<GroupKey, Integer>> processor = new AbstractParallelGroupedDataStoringAggregationProcessor<Integer, Integer>(ConcurrencyTestsUtil.getSharedExecutor(), receivers, "Sum") {
@Override
protected void storeElement(GroupedDataEntry<Integer> element) {
if (element.getDataEntry() < 0) {
@@ -16,7 +16,7 @@ public class TestParallelAggregationProcessors extends AbstractTestParallelAvera
@Test
public void testSumAggregationProcessor() throws InterruptedException {
Processor<GroupedDataEntry<Number>, Map<GroupKey, Number>> sumAggregationProcessor = ParallelGroupedNumberDataSumAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getExecutor(), receivers);
Processor<GroupedDataEntry<Number>, Map<GroupKey, Number>> sumAggregationProcessor = ParallelGroupedNumberDataSumAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getSharedExecutor(), receivers);
Collection<GroupedDataEntry<Number>> elements = createElements();
ConcurrencyTestsUtil.processElements(sumAggregationProcessor, elements);
@@ -27,7 +27,7 @@ public class TestParallelAggregationProcessors extends AbstractTestParallelAvera
@Test
public void testMedianAggregationProcessor() throws InterruptedException {
Processor<GroupedDataEntry<Number>, Map<GroupKey, Number>> medianAggregationProcessor = ParallelGroupedNumberDataMedianAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getExecutor(), receivers);
Processor<GroupedDataEntry<Number>, Map<GroupKey, Number>> medianAggregationProcessor = ParallelGroupedNumberDataMedianAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getSharedExecutor(), receivers);
Collection<GroupedDataEntry<Number>> elements = createElements();
ConcurrencyTestsUtil.processElements(medianAggregationProcessor, elements);
@@ -38,7 +38,7 @@ public class TestParallelAggregationProcessors extends AbstractTestParallelAvera
@Test
public void testMaxAggregationProcessor() throws InterruptedException {
Processor<GroupedDataEntry<Number>, Map<GroupKey, Number>> maxAggregationProcessor = ParallelGroupedNumberDataMaxAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getExecutor(), receivers);
Processor<GroupedDataEntry<Number>, Map<GroupKey, Number>> maxAggregationProcessor = ParallelGroupedNumberDataMaxAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getSharedExecutor(), receivers);
Collection<GroupedDataEntry<Number>> elements = createElements();
ConcurrencyTestsUtil.processElements(maxAggregationProcessor, elements);
@@ -49,7 +49,7 @@ public class TestParallelAggregationProcessors extends AbstractTestParallelAvera
@Test
public void testMinAggregationProcessor() throws InterruptedException {
Processor<GroupedDataEntry<Number>, Map<GroupKey, Number>> minAggregationProcessor = ParallelGroupedNumberDataMinAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getExecutor(), receivers);
Processor<GroupedDataEntry<Number>, Map<GroupKey, Number>> minAggregationProcessor = ParallelGroupedNumberDataMinAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getSharedExecutor(), receivers);
Collection<GroupedDataEntry<Number>> elements = createElements();
ConcurrencyTestsUtil.processElements(minAggregationProcessor, elements);
@@ -60,7 +60,7 @@ public class TestParallelAggregationProcessors extends AbstractTestParallelAvera
@Test
public void testCountAggregationProcessor() throws InterruptedException {
Processor<GroupedDataEntry<Object>, Map<GroupKey, Number>> countAggregationProcessor = ParallelGroupedDataCountAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getExecutor(), receivers);
Processor<GroupedDataEntry<Object>, Map<GroupKey, Number>> countAggregationProcessor = ParallelGroupedDataCountAggregationProcessor.getDefinition().construct(ConcurrencyTestsUtil.getSharedExecutor(), receivers);
@SuppressWarnings("unchecked")
Collection<GroupedDataEntry<Object>> elements = (Collection<GroupedDataEntry<Object>>)(Collection<?>) createElements();
ConcurrencyTestsUtil.processElements(countAggregationProcessor, elements);
@@ -19,7 +19,7 @@ public class TestParallelAveragingProcessors extends AbstractTestParallelAveragi
@Test
public void testAverageAggregationProcessor() throws InterruptedException {
Processor<GroupedDataEntry<Number>, Map<GroupKey, AverageWithStats<Number>>> averageAggregationProcessor = ParallelGroupedNumberDataAverageAggregationProcessor
.getDefinition().construct(ConcurrencyTestsUtil.getExecutor(), receivers);
.getDefinition().construct(ConcurrencyTestsUtil.getSharedExecutor(), receivers);
Collection<GroupedDataEntry<Number>> elements = createElements();
ConcurrencyTestsUtil.processElements(averageAggregationProcessor, elements);
averageAggregationProcessor.finish();
@@ -37,7 +37,7 @@ public class TestParallelGroupedDataCollectingAsSetProcessor {
@Test
public void testDataCollecting() throws InterruptedException {
Processor<GroupedDataEntry<Double>, Map<GroupKey, HashSet<Double>>> collectingProcessor = new ParallelGroupedDataCollectingAsSetProcessor<Double>(ConcurrencyTestsUtil.getExecutor(), receivers);
Processor<GroupedDataEntry<Double>, Map<GroupKey, HashSet<Double>>> collectingProcessor = new ParallelGroupedDataCollectingAsSetProcessor<Double>(ConcurrencyTestsUtil.getSharedExecutor(), receivers);
Collection<GroupedDataEntry<Double>> elements = createElements();
ConcurrencyTestsUtil.processElements(collectingProcessor, elements);
@@ -63,7 +63,7 @@ public class TestFunctionManagerAsFunctionRegistry {
DataRetrieverChainDefinitionRegistry dataRetrieverChainDefinitionRegistry = new DataRetrieverChainDefinitionManager();
AggregationProcessorDefinitionRegistry aggregationProcessorDefinitionRegistry = new AggregationProcessorDefinitionManager();
QueryDefinitionDTORegistry queryDefinitionRegistry = new QueryDefinitionDTOManager();
ModifiableDataMiningServer server = new DataMiningServerImpl(ConcurrencyTestsUtil.getExecutor(), functionManager,
ModifiableDataMiningServer server = new DataMiningServerImpl(ConcurrencyTestsUtil.getSharedExecutor(), functionManager,
dataSourceProviderRegistry,
dataRetrieverChainDefinitionRegistry,
aggregationProcessorDefinitionRegistry,
@@ -22,7 +22,7 @@ import com.sap.sse.datamining.test.domain.impl.Test_TeamImpl;
public final class ComponentTestsUtil {
private final static ProcessorFactory processorFactory = new ProcessorFactory(ConcurrencyTestsUtil.getExecutor());
private final static ProcessorFactory processorFactory = new ProcessorFactory(ConcurrencyTestsUtil.getSharedExecutor());
public static ProcessorFactory getProcessorFactory() {
return processorFactory;
@@ -19,7 +19,7 @@ public class ConcurrencyTestsUtil extends TestsUtil {
private static final int THREAD_POOL_SIZE = Math.max(Runtime.getRuntime().availableProcessors(), 3);
private static final ExecutorService executor = new DataMiningExecutorService(THREAD_POOL_SIZE);
public static ExecutorService getExecutor() {
public static ExecutorService getSharedExecutor() {
return executor;
}
@@ -1,6 +1,7 @@
package com.sap.sse.datamining.test.util;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.ExecutorService;
import com.sap.sse.datamining.ModifiableDataMiningServer;
import com.sap.sse.datamining.components.management.AggregationProcessorDefinitionRegistry;
@@ -55,12 +56,16 @@ public class TestsUtil {
}
public static ModifiableDataMiningServer createNewServer() {
return createNewServer(ConcurrencyTestsUtil.getSharedExecutor());
}
public static ModifiableDataMiningServer createNewServer(ExecutorService executor) {
FunctionRegistry functionRegistry = new FunctionManager();
DataSourceProviderRegistry dataSourceProviderRegistry = new DataSourceProviderManager();
DataRetrieverChainDefinitionRegistry dataRetrieverChainDefinitionRegistry = new DataRetrieverChainDefinitionManager();
AggregationProcessorDefinitionRegistry aggregationProcessorDefinitionRegistry = new AggregationProcessorDefinitionManager();
QueryDefinitionDTORegistry queryDefinitionRegistry = new QueryDefinitionDTOManager();
return new DataMiningServerImpl(ConcurrencyTestsUtil.getExecutor(), functionRegistry,
return new DataMiningServerImpl(executor, functionRegistry,
dataSourceProviderRegistry,
dataRetrieverChainDefinitionRegistry,
aggregationProcessorDefinitionRegistry,
@@ -2,18 +2,29 @@ package com.sap.sse.datamining.test.util.components;
import com.sap.sse.datamining.components.ProcessorInstructionHandler;
import com.sap.sse.datamining.impl.components.AbstractProcessorInstruction;
import com.sap.sse.datamining.impl.components.ProcessorInstructionPriority;
public class StatefulBlockingInstruction extends AbstractProcessorInstruction<Object> {
public class StatefulBlockingInstruction<ResultType> extends AbstractProcessorInstruction<ResultType> {
private final int duration;
protected final long blockDuration;
protected final ResultType result;
private boolean runWasCalled;
private boolean computeResultWasCalled;
private boolean computeResultWasFinished;
public StatefulBlockingInstruction(ProcessorInstructionHandler<Object> handler, int duration) {
super(handler);
this.duration = duration;
public StatefulBlockingInstruction(ProcessorInstructionHandler<ResultType> handler, long blockDuration) {
this(handler, 0, blockDuration, null);
}
public StatefulBlockingInstruction(ProcessorInstructionHandler<ResultType> handler, ProcessorInstructionPriority priority, long blockDuration, ResultType result) {
this(handler, priority.asIntValue(), blockDuration, result);
}
public StatefulBlockingInstruction(ProcessorInstructionHandler<ResultType> handler, int priority, long blockDuration, ResultType result) {
super(handler, priority);
this.blockDuration = blockDuration;
this.result = result;
}
@Override
@@ -23,13 +34,22 @@ public class StatefulBlockingInstruction extends AbstractProcessorInstruction<Ob
}
@Override
protected Object computeResult() throws Exception {
protected ResultType computeResult() throws Exception {
computeResultWasCalled = true;
if (duration > 0) {
Thread.sleep(duration);
actionBeforeBlock();
if (blockDuration > 0) {
Thread.sleep(blockDuration);
}
actionAfterBlock();
computeResultWasFinished = true;
return null;
return result;
}
protected void actionBeforeBlock() { }
protected void actionAfterBlock() { }
public long getBlockDuration() {
return blockDuration;
}
public boolean runWasCalled() {