Implemented a timeouting mechanism for queries with an issue

The problem is, that the processing still continous, because there is no
abort mechanism for processors right now. The query just stops waiting for
the result and throws a TimeoutException.
This commit is contained in:
Lennart Hensler committed 2014-03-02 11:47:55 +01:00
1 parent a0f6af0318
commit f00dc963ba
3 files changed
+105 -27

No files matched your search

@@ -2,10 +2,13 @@ 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.Arrays;
import java.util.Collection;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
@@ -13,7 +16,6 @@ import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import org.junit.Ignore;
import org.junit.Test;
import com.sap.sse.datamining.Query;
@@ -32,12 +34,15 @@ import com.sap.sse.datamining.test.util.FunctionTestsUtil;
public class TestProcessorQuery {
@SuppressWarnings("unused")
private boolean receivedElementOrFinished;
@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);
verifyResult(queryWithStandardWorkflow.run(), expectedResult);
}
private Collection<Number> createDataSource() {
@@ -80,7 +85,7 @@ public class TestProcessorQuery {
*/
private Query<Double> createQueryWithStandardWorkflow(Collection<Number> dataSource) {
ThreadPoolExecutor executor = ConcurrencyTestsUtil.getExecutor();
ProcessorQuery<Double, Iterable<Number>> query = new ProcessorQuery<Double, Iterable<Number>>(dataSource);
ProcessorQuery<Double, Iterable<Number>> query = new ProcessorQuery<Double, Iterable<Number>>(executor, dataSource);
Collection<Processor<Map<GroupKey, Double>>> aggregationResultReceivers = asCollection(query.getResultReceiver());
Processor<GroupedDataEntry<Double>> sumAggregator =
@@ -117,7 +122,7 @@ public class TestProcessorQuery {
return result;
}
private void verifyResultWith(QueryResult<Double> result, QueryResult<Double> expectedResult) {
private void verifyResult(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()));
@@ -126,22 +131,74 @@ public class TestProcessorQuery {
// assertThat("Value decimals aren't correct.", result.getValueDecimals(), is(expectedResult.getValueDecimals()));
}
@Ignore
@Test
@Test(timeout=2000)
public void testQueryTimeouting() throws TimeoutException {
ProcessorQuery<Double, Iterable<Number>> query = new ProcessorQuery<Double, Iterable<Number>>(createDataSource());
query.setFirstProcessor(createBlockingProcessor());
query.run(5000, TimeUnit.MILLISECONDS);
ProcessorQuery<Double, Iterable<Number>> query = new ProcessorQuery<Double, Iterable<Number>>(ConcurrencyTestsUtil.getExecutor(), createDataSource());
Processor<Double> resultReceiver = new Processor<Double>() {
@Override
public void onElement(Double element) {
receivedElementOrFinished = true;
}
@Override
public void finish() throws InterruptedException {
receivedElementOrFinished = true;
}
};
query.setFirstProcessor(createBlockingProcessor(1000, resultReceiver));
try {
query.run(500, TimeUnit.MILLISECONDS);
fail("The previous line should throw a timeout exception");
} catch (TimeoutException e) {
// A timeout exception is expected
}
// ConcurrencyTestsUtil.sleepFor(1000); // Wait if a result is received
// assertThat("The processing should be aborted", receivedElementOrFinished, is(false));
}
private Processor<Iterable<Number>> createBlockingProcessor() {
return new AbstractSimpleParallelProcessor<Iterable<Number>, Double>(ConcurrencyTestsUtil.getExecutor(), new ArrayList<Processor<Double>>()) {
private Processor<Iterable<Number>> createBlockingProcessor(final long timeToBlockInMillis, Processor<Double> resultReceiver) {
return new AbstractSimpleParallelProcessor<Iterable<Number>, Double>(ConcurrencyTestsUtil.getExecutor(), Arrays.asList(resultReceiver)) {
@Override
protected Callable<Double> createInstruction(Iterable<Number> element) {
return new Callable<Double>() {
@Override
public Double call() throws Exception {
while (true) { }
Thread.sleep(timeToBlockInMillis);
return 0.0;
}
};
}
};
}
@Test
public void testQueryWithTimeoutAndNonBlockingProcess() throws TimeoutException {
ProcessorQuery<Double, Iterable<Number>> query = new ProcessorQuery<Double, Iterable<Number>>(ConcurrencyTestsUtil.getExecutor(), createDataSource());
String keyValue = "Sum";
query.setFirstProcessor(createSumBuildingProcessor(query, keyValue));
Map<GroupKey, Double> expectedResult = new HashMap<>();
expectedResult.put(new GenericGroupKey<String>(keyValue), 10358.0);
assertThat(query.run(500, TimeUnit.MILLISECONDS).getResults(), is(expectedResult));
}
private AbstractSimpleParallelProcessor<Iterable<Number>, Map<GroupKey, Double>> createSumBuildingProcessor(
ProcessorQuery<Double, Iterable<Number>> query, final String keyValue) {
return new AbstractSimpleParallelProcessor<Iterable<Number>, Map<GroupKey, Double>>(ConcurrencyTestsUtil.getExecutor(),
Arrays.asList(query.getResultReceiver())) {
@Override
protected Callable<Map<GroupKey, Double>> createInstruction(final Iterable<Number> element) {
return new Callable<Map<GroupKey,Double>>() {
@Override
public Map<GroupKey, Double> call() throws Exception {
Map<GroupKey, Double> result = new HashMap<>();
double sum = 0;
for (Number number : element) {
sum += number.getValue();
}
result.put(new GenericGroupKey<String>(keyValue), sum);
return result;
}
};
}
@@ -15,6 +15,10 @@ public class Number {
this.value = value;
}
public int getValue() {
return value;
}
@Dimension("length")
public int getLength() {
return String.valueOf(value).length();
@@ -1,7 +1,10 @@
package com.sap.sse.datamining.impl.components;
import java.util.Map;
import java.util.Timer;
import java.util.TimerTask;
import java.util.Map.Entry;
import java.util.concurrent.Executor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.logging.Level;
@@ -20,14 +23,16 @@ public class ProcessorQuery<AggregatedType, DataSourceType> implements Query<Agg
private final DataSourceType dataSource;
private Processor<DataSourceType> firstProcessor;
private Executor executor;
private final ProcessResultReceiver resultReceiver;
private final Object monitorObject = new Object();
private boolean workIsDone = false;
private boolean processorTimedOut = true;
public ProcessorQuery(DataSourceType dataSource) {
public ProcessorQuery(Executor executor, DataSourceType dataSource) {
this.executor = executor;
this.dataSource = dataSource;
resultReceiver = new ProcessResultReceiver();
}
@@ -60,30 +65,28 @@ public class ProcessorQuery<AggregatedType, DataSourceType> implements Query<Agg
}
private QueryResult<AggregatedType> processQuery(long timeoutInMillis) throws InterruptedException, TimeoutException {
firstProcessor.onElement(dataSource);
firstProcessor.finish();
setUpTimeout(timeoutInMillis);
waitTillWorkIsDone();
return resultReceiver.getResult();
}
private void setUpTimeout(long timeoutInMillis) {
Thread timeoutThread = new Thread(new Runnable() {
executor.execute(new Runnable() {
@Override
public void run() {
synchronized (monitorObject) {
monitorObject.notify();
try {
firstProcessor.onElement(dataSource);
firstProcessor.finish();
} catch (InterruptedException e) {
LOGGER.log(Level.WARNING, "The query processing got interrupted.", e);
}
}
});
timeoutThread.start();
waitTillWorkIsDone(timeoutInMillis);
return resultReceiver.getResult();
}
private void waitTillWorkIsDone() throws InterruptedException, TimeoutException {
private void waitTillWorkIsDone(long timeoutInMillis) throws InterruptedException, TimeoutException {
setUpTimeoutTimer(timeoutInMillis);
synchronized (monitorObject) {
while (!workIsDone) {
monitorObject.wait();
if (processorTimedOut) {
// TODO Abort the processing
throw new TimeoutException("The query processing timed out");
}
}
@@ -91,6 +94,20 @@ public class ProcessorQuery<AggregatedType, DataSourceType> implements Query<Agg
processorTimedOut = true;
}
}
private void setUpTimeoutTimer(long timeoutInMillis) {
if (timeoutInMillis > 0) {
Timer timeoutTimer = new Timer();
timeoutTimer.schedule(new TimerTask() {
@Override
public void run() {
synchronized (monitorObject) {
monitorObject.notify();
}
}
}, timeoutInMillis);
}
}
Processor<Map<GroupKey, AggregatedType>> getResultReceiver() {
return resultReceiver;