Made ManeuverAndWindImporter to operate in parallel

This commit is contained in:
Vladislav Chumak
2019-02-18 15:32:19 +01:00
parent 019a9b7a78
commit a95994663e
@@ -13,6 +13,9 @@ import java.time.temporal.ChronoUnit;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import org.apache.http.Header;
import org.apache.http.HttpResponse;
@@ -40,11 +43,11 @@ import com.sap.sailing.windestimation.data.persistence.maneuver.RaceWithManeuver
import com.sap.sailing.windestimation.data.persistence.maneuver.RaceWithManeuverForEstimationPersistenceManager;
import com.sap.sailing.windestimation.data.persistence.twdtransition.RaceWithWindSourcesPersistenceManager;
import com.sap.sailing.windestimation.data.serialization.CompetitorTrackWithEstimationDataJsonDeserializer;
import com.sap.sailing.windestimation.data.serialization.ManeuverForDataAnalysisJsonSerializer;
import com.sap.sailing.windestimation.data.serialization.LabeledManeuverForEstimationJsonSerializer;
import com.sap.sailing.windestimation.data.serialization.ManeuverForDataAnalysisJsonSerializer;
import com.sap.sailing.windestimation.data.transformer.CompetitorTrackTransformer;
import com.sap.sailing.windestimation.data.transformer.CompleteManeuverCurveWithEstimationDataToManeuverForDataAnalysisTransformer;
import com.sap.sailing.windestimation.data.transformer.CompleteManeuverCurveWithEstimationDataToLabelledManeuverForEstimationTransformer;
import com.sap.sailing.windestimation.data.transformer.CompleteManeuverCurveWithEstimationDataToManeuverForDataAnalysisTransformer;
import com.sap.sailing.windestimation.util.LoggingUtil;
/**
@@ -70,6 +73,8 @@ public class ManeuverAndWindImporter {
private final ManeuverForDataAnalysisJsonSerializer maneuverForDataAnalysisJsonSerializer;
private final LabeledManeuverForEstimationJsonSerializer maneuverForEstimationJsonSerializer;
private boolean skipRace;
private static final int NUMBER_OF_THREADS = 50; //high number due to HTTP requests
private final ExecutorService executorService = Executors.newFixedThreadPool(NUMBER_OF_THREADS);
public ManeuverAndWindImporter() throws UnknownHostException {
this.completeManeuverCurvePersistanceManager = new RaceWithCompleteManeuverCurvePersistenceManager();
@@ -94,8 +99,8 @@ public class ManeuverAndWindImporter {
importer.importAllRegattas();
}
public void importAllRegattas()
throws IllegalStateException, ClientProtocolException, IOException, ParseException, URISyntaxException {
public void importAllRegattas() throws IllegalStateException, ClientProtocolException, IOException, ParseException,
URISyntaxException, InterruptedException {
skipRace = startFromRegattaName != null;
LoggingUtil.logInfo("Importer for CompleteManeuverCurveWithEstimationData just started");
LoggingUtil.logInfo("Dropping old database");
@@ -118,9 +123,17 @@ public class ManeuverAndWindImporter {
+ Math.round(100.0 * i / numberOfRegattas) + "%): \"" + regattaName + "\"");
importRegatta(regattaName, importStatistics);
}
LoggingUtil.logInfo("Import finished");
importStatistics.regattasCount = regattasJson.size();
logImportStatistics(importStatistics);
boolean success = executorService.awaitTermination(24, TimeUnit.HOURS);
executorService.shutdown();
if (success) {
LoggingUtil.logInfo("Import finished");
synchronized (importStatistics) {
importStatistics.regattasCount = regattasJson.size();
logImportStatistics(importStatistics);
}
} else {
LoggingUtil.logInfo("Process was terminated");
}
}
private void logImportStatistics(ImportStatistics importStatistics) {
@@ -145,7 +158,9 @@ public class ManeuverAndWindImporter {
try {
regattaJson = (JSONObject) getHttpResponseAsJson(regattaName, null, getRegatta);
} catch (Exception e) {
importStatistics.ignoredRegattas += 1;
synchronized (importStatistics) {
importStatistics.ignoredRegattas++;
}
LoggingUtil.logInfo("Error while processing regatta: " + regattaName);
return;
}
@@ -161,17 +176,25 @@ public class ManeuverAndWindImporter {
if ((boolean) race.get("isTracked") && !(boolean) race.get("isLive")
&& (boolean) race.get("hasGpsData") && (boolean) race.get("hasWindData")) {
String trackedRaceName = (String) race.get("trackedRaceName");
LoggingUtil.logInfo("Processing race nr. " + ++i + ": \"" + trackedRaceName + "\"");
try {
importRace(regattaName, trackedRaceName, importStatistics);
} catch (Exception e) {
importStatistics.ingoredRaces += 1;
LoggingUtil
.logInfo("Error while processing race nr. " + i + ": \"" + trackedRaceName + "\"");
}
i++;
final int raceNumber = i;
LoggingUtil.logInfo("Processing race nr. " + raceNumber + ": \"" + trackedRaceName + "\"");
executorService.execute(() -> {
try {
importRace(regattaName, trackedRaceName, importStatistics);
} catch (Exception e) {
synchronized (importStatistics) {
importStatistics.ingoredRaces += 1;
}
LoggingUtil.logInfo("Error while processing race nr. " + raceNumber + ": \""
+ trackedRaceName + "\"");
}
});
}
}
importStatistics.racesCount += racesJson.size();
synchronized (importStatistics) {
importStatistics.racesCount += racesJson.size();
}
}
}
}
@@ -217,7 +240,9 @@ public class ManeuverAndWindImporter {
}
LoggingUtil.logInfo(
"Imported " + windFixesCount + " wind fixes from " + windSourcesJson.size() + " wind sources");
importStatistics.racesWithHighQualityWindData++;
synchronized (importStatistics) {
importStatistics.racesWithHighQualityWindData++;
}
} else {
LoggingUtil.logInfo("No high quality wind fixes contained");
}
@@ -274,8 +299,10 @@ public class ManeuverAndWindImporter {
}
LoggingUtil.logInfo(
"Imported " + competitorTracks.size() + " competitor tracks with " + maneuversCount + " maneuvers");
importStatistics.competitorTracksCount += competitorTracks.size();
importStatistics.maneuversCount += maneuversCount;
synchronized (importStatistics) {
importStatistics.competitorTracksCount += competitorTracks.size();
importStatistics.maneuversCount += maneuversCount;
}
}
private JSONObject getHttpResponseAsJson(String trackedRegattaName, String trackedRaceName,