Commit ee6980a0 authored by Luna Riegel's avatar Luna Riegel
Browse files

Refactor: Add one-time retrying of checks

Checker will now retry to check Features that were missed in the first
checking attempt due to being unable to load the Features from the
database
parent 7cafd309
......@@ -90,7 +90,6 @@ import java.util.Map.Entry;
import java.util.Set;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Stream;
/**
* The main container class for checking. It contains the logic for validation,
......@@ -650,27 +649,71 @@ public class Checker {
private void checkCityModel(CityDoctorModel model, ProgressListener l){
List<GmlId> features = model.getFeatureIds();
CityObjectCache cache = model.getCache();
float featureSum = model.getNumberOfFeatures();
AtomicInteger checkedCount = new AtomicInteger(0);
AtomicInteger missingCount = new AtomicInteger(0);
// clear global errors
model.getGlobalErrors().clear();
CityObjectCache cache = model.getCache();
logger.trace("Setting up ThreadPool");
int threadCount = Runtime.getRuntime().availableProcessors() * 5;
ExecutorService exec = new ThreadPoolExecutor
(threadCount, threadCount, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>());
try{
long startTime = System.nanoTime();
List<Future<GmlId>> futures = runChecksOnFeatures(exec, cache, features, checkedCount, l);
List<GmlId> missedFeatures = getMissedFeatures(futures);
if (!missedFeatures.isEmpty()){
logger.warn("Checker could not load some CityObjects from database!");
logger.info("Waiting a second before retrying...");
Thread.sleep(1000);
futures = runChecksOnFeatures(exec, cache, missedFeatures, checkedCount, l);
missedFeatures = getMissedFeatures(futures);
if (!missedFeatures.isEmpty()){
//TODO: Localize String
logger.error("Checker still unable to load some CityObjects from database!");
logger.debug(missedFeatures.toString());
}
}
long endTime = System.nanoTime();
if (logger.isDebugEnabled()){
logger.debug("Checks finished in {} ms", (endTime - startTime) / 1_000_000);
}
} catch (InterruptedException e){
// No interrupts implemented, so this block should actually never be entered
//TODO: Localize String
logger.error("Validation aborted: Checker got interrupted", e);
Thread.currentThread().interrupt();
} finally {
exec.shutdown();
try {
if(!exec.awaitTermination(1, TimeUnit.MINUTES)){
exec.shutdownNow();
}
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
}
}
private List<Future<GmlId>> runChecksOnFeatures(ExecutorService exec, CityObjectCache cache, List<GmlId> ids, AtomicInteger checkedCount,
ProgressListener l)
throws InterruptedException {
float featureSum = ids.size() + (float) checkedCount.get();
long startTime = System.nanoTime();
for (GmlId id : features) {
exec.execute(()->{
logger.trace("Queueing up Checker tasks");
List<Callable<GmlId>> tasks = new ArrayList<>();
for (GmlId id : ids) {
tasks.add(()->{
CityObject co = cache.get(id);
if (co == null) {
missingCount.incrementAndGet();
if (l!=null){
l.updateProgress((checkedCount.get() + missingCount.get()) / featureSum);
}
return;
return id;
}
if(Thread.interrupted()){
Thread.currentThread().interrupt();
return null;
}
co.prepareForChecking();
executeChecksForCityObject(co);
......@@ -678,63 +721,28 @@ public class Checker {
cache.put(co);
checkedCount.incrementAndGet();
if (l!=null){
l.updateProgress((checkedCount.get() + missingCount.get()) / featureSum);
l.updateProgress((checkedCount.get()) / featureSum);
}
return null;
});
}
exec.shutdown();
try {
exec.awaitTermination(30, TimeUnit.DAYS);
long endTime = System.nanoTime();
if (logger.isDebugEnabled()){
logger.debug("Checks finished in {} ms", (endTime - startTime) / 1_000_000);
}
if (missingCount.get() > 0) {
int missing = missingCount.get();
//TODO: Localize String
String msg = String.format("Checker could not load %d/%.0f features (%.2f%%)", missing, featureSum, (missing/featureSum) * 100);
logger.warn(msg);
}
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
return exec.invokeAll(tasks);
}
private void checkCityModelOLD(CityDoctorModel model, ProgressListener l) {
Stream<CityObject> features = model.createFeatureStream();
float featureSum = model.getNumberOfFeatures();
boolean lowMemoryMode = config.getParserConfiguration().useLowMemoryConsumption();
// clear global errors
model.getGlobalErrors().clear();
long startTime = System.nanoTime();
// stupid lamda with final variable restrictions
int[] currentFeature = new int[1];
features.forEach(co -> {
if (lowMemoryMode) {
// no edges have been created yet, create them
co.prepareForChecking();
}
// check every feature
executeChecksForCityObject(co);
if (lowMemoryMode) {
// low memory consumption, remove edges again
co.clearMetaInformation();
}
if (l != null) {
currentFeature[0]++;
l.updateProgress(currentFeature[0] / featureSum);
private List<GmlId> getMissedFeatures(List<Future<GmlId>> futures) throws InterruptedException{
List<GmlId> missedList = new ArrayList<>();
for (Future<GmlId> future : futures) {
try{
GmlId gmlId = future.get();
if (gmlId != null) {
missedList.add(gmlId);
}
} catch (ExecutionException e){
logger.debug("A Task failed due to an unexpected exception", e);
logger.debug(e.getCause());
}
});
long endTime = System.nanoTime();
model.preCacheShownFeatures();
logger.info("Checks finished in " + (endTime -startTime) / 1_000_000 + " ms");
}
return missedList;
}
private boolean filterObject(CityObject co) {
......
Supports Markdown
0% or .
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment