mirror of
https://github.com/eclipse-sailing-analytics/sailing-analytics.git
synced 2026-10-08 13:20:57 +00:00
Cleaned the instruction wrapping of the abstract parallel processor
This commit is contained in:
1 parent
f00dc963ba
commit
a7d837f2d6
3 files changed
+34
-25
No files matched your search
+2
-16
@@ -28,6 +28,7 @@ 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.BlockingProcessor;
|
||||
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;
|
||||
@@ -144,7 +145,7 @@ public class TestProcessorQuery {
|
||||
receivedElementOrFinished = true;
|
||||
}
|
||||
};
|
||||
query.setFirstProcessor(createBlockingProcessor(1000, resultReceiver));
|
||||
query.setFirstProcessor(new BlockingProcessor<Iterable<Number>, Double>(ConcurrencyTestsUtil.getExecutor(), Arrays.asList(resultReceiver), (long) 1000));
|
||||
|
||||
try {
|
||||
query.run(500, TimeUnit.MILLISECONDS);
|
||||
@@ -157,21 +158,6 @@ public class TestProcessorQuery {
|
||||
// assertThat("The processing should be aborted", receivedElementOrFinished, is(false));
|
||||
}
|
||||
|
||||
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 {
|
||||
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());
|
||||
|
||||
+28
@@ -0,0 +1,28 @@
|
||||
package com.sap.sse.datamining.test.components.util;
|
||||
|
||||
import java.util.Collection;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.Executor;
|
||||
|
||||
import com.sap.sse.datamining.components.Processor;
|
||||
import com.sap.sse.datamining.impl.components.AbstractSimpleParallelProcessor;
|
||||
|
||||
public class BlockingProcessor<InputType, ResultType> extends AbstractSimpleParallelProcessor<InputType, ResultType> {
|
||||
private final long timeToBlockInMillis;
|
||||
|
||||
public BlockingProcessor(Executor executor, Collection<Processor<ResultType>> resultReceivers, long timeToBlockInMillis) {
|
||||
super(executor, resultReceivers);
|
||||
this.timeToBlockInMillis = timeToBlockInMillis;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected Callable<ResultType> createInstruction(InputType element) {
|
||||
return new Callable<ResultType>() {
|
||||
@Override
|
||||
public ResultType call() throws Exception {
|
||||
Thread.sleep(timeToBlockInMillis);
|
||||
return null;
|
||||
}
|
||||
};
|
||||
}
|
||||
}
|
||||
+4
-9
@@ -4,10 +4,7 @@ import java.util.Collection;
|
||||
import java.util.HashSet;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.Executor;
|
||||
import java.util.concurrent.FutureTask;
|
||||
import java.util.concurrent.RunnableFuture;
|
||||
import java.util.concurrent.locks.ReentrantReadWriteLock;
|
||||
import java.util.logging.Level;
|
||||
import java.util.logging.Logger;
|
||||
@@ -33,20 +30,18 @@ public abstract class AbstractPartitioningParallelProcessor<InputType, WorkingTy
|
||||
@Override
|
||||
public void onElement(InputType element) {
|
||||
for (WorkingType partialElement : partitionElement(element)) {
|
||||
Callable<ResultType> instruction = createInstruction(partialElement);
|
||||
final Callable<ResultType> instruction = createInstruction(partialElement);
|
||||
if (isInstructionValid(instruction)) {
|
||||
final RunnableFuture<ResultType> runnableInstruction = new FutureTask<>(instruction);
|
||||
Runnable instructionWrapper = new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
runnableInstruction.run();
|
||||
try {
|
||||
ResultType result = runnableInstruction.get();
|
||||
ResultType result = instruction.call();
|
||||
if (isResultValid(result)) {
|
||||
forwardResultToReceivers(result);
|
||||
}
|
||||
} catch (InterruptedException | ExecutionException e) {
|
||||
LOGGER.log(Level.FINEST, "Error getting the result from the instruction: ", e);
|
||||
} catch (Exception e) {
|
||||
LOGGER.log(Level.FINEST, "An error occured during the processing of an instruction: ", e);
|
||||
} finally {
|
||||
AbstractPartitioningParallelProcessor.this.unfinishedInstructionsCounter.decrement();
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user