Introduced consolidation of bravo fixes

This commit is contained in:
Papick Garcia Taboada
2017-02-15 17:38:36 +01:00
parent 3ebf92f1cf
commit 69463e33a1
17 changed files with 581 additions and 122 deletions
@@ -9,15 +9,31 @@ import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import com.sap.sailing.domain.common.sensordata.BravoSensorDataMetadata;
import com.sap.sailing.domain.common.tracking.DoubleVectorFix;
import com.sap.sailing.domain.common.tracking.impl.BravoFixImpl;
import com.sap.sailing.domain.trackfiles.TrackFileImportDeviceIdentifier;
import com.sap.sailing.domain.trackimport.DoubleVectorFixImporter;
import com.sap.sailing.domain.trackimport.FormatNotSupportedException;
import com.sap.sailing.server.trackfiles.impl.BravoDataImporterImpl;
import com.sap.sailing.server.trackfiles.impl.doublefix.DownsamplerTo1HzProcessor;
import com.sap.sailing.server.trackfiles.impl.doublefix.LearningBatchProcessor;
public class BravoDataImportTest {
private final DoubleVectorFixImporter bravoDataImporter = new BravoDataImporterImpl();
private LearningBatchProcessor batchProcessor;
private DownsamplerTo1HzProcessor downsampler;
private final DoubleVectorFixImporter bravoDataImporter = new BravoDataImporterImpl() {
protected com.sap.sailing.server.trackfiles.impl.doublefix.DoubleFixProcessor createProcessor(
BravoSensorDataMetadata metadata,
DoubleVectorFixImporter.Callback callback,
TrackFileImportDeviceIdentifier trackIdentifier) {
batchProcessor = new LearningBatchProcessor(5000, 5000, callback, trackIdentifier);
downsampler = new DownsamplerTo1HzProcessor(metadata.trackColumnCount, batchProcessor);
return downsampler;
};
};
private int callbackCallCount = 0;
private double sumRideHeightInMeters = 0.0;
@@ -48,39 +64,40 @@ public class BravoDataImportTest {
}
private void testImport(ImportData importData) throws FormatNotSupportedException, IOException {
bravoDataImporter.importFixes(importData.getInputStream(), (fixes, device) -> {
for (DoubleVectorFix fix : fixes) {
callbackCallCount++;
sumRideHeightInMeters += new BravoFixImpl(fix).getRideHeight().getMeters();
}
}, "filename", "source");
Assert.assertEquals(importData.expectedFixesCount, callbackCallCount);
Assert.assertEquals(importData.expectedAverageRideHeight,
callbackCallCount == 0 ? sumRideHeightInMeters : sumRideHeightInMeters / callbackCallCount, 0.00001);
Assert.assertEquals(importData.expectedFixesCount, downsampler.getCountSourceTtl());
Assert.assertEquals(importData.expectedFixesConsolidated, downsampler.getCountImportedTtl());
Assert.assertEquals(importData.expectedFixesConsolidated, callbackCallCount);
}
private enum ImportData {
// find out the number of fixes using the following bash line:
// tail -n +5 Undefined\ Race\ -\ BRAVO.txt | awk '{if ($5<$6) print v=$5; else v=$6; sum+=v; count++; } END {print "Sum: ", sum, "Count: ", count, "Average:", sum/count;}'
FILE_UNDEFINED_RACE_BRAVO(870, 0.661037) {
FILE_UNDEFINED_RACE_BRAVO(870, 89) {
@Override
protected InputStream getInputStream() {
return getClass().getResourceAsStream("/Undefined Race - BRAVO.txt");
}
},
DUMMY_DEFAULT_HEADER_NO_DATA(0, 0.0) {
DUMMY_DEFAULT_HEADER_NO_DATA(0, 0) {
@Override
protected InputStream getInputStream() {
return new ByteArrayInputStream(HEADER_ORDER_DEFAULT.getBytes(StandardCharsets.UTF_8));
}
},
DUMMY_DEAFULT_HEADER_ONE_LINE(1, 0.661044) {
DUMMY_DEAFULT_HEADER_ONE_LINE(1, 1) {
@Override
protected InputStream getInputStream() {
return new ByteArrayInputStream((HEADER_ORDER_DEFAULT + DUMMY_CONTENT).getBytes(StandardCharsets.UTF_8));
}
},
DUMMY_RANDOM_HAEDER_ONE_LINE(1, 0.092727) {
DUMMY_RANDOM_HAEDER_ONE_LINE(1, 1) {
@Override
protected InputStream getInputStream() {
return new ByteArrayInputStream((HEADER_ORDER_RANDOM + DUMMY_CONTENT).getBytes(StandardCharsets.UTF_8));
@@ -88,11 +105,11 @@ public class BravoDataImportTest {
};
private final int expectedFixesCount;
private final double expectedAverageRideHeight;
private final int expectedFixesConsolidated;
private ImportData(int expectedFixesCount, double expectedAverageRideHeight) {
private ImportData(int expectedFixesCount, int expectedFixesConsolidated) {
this.expectedFixesCount = expectedFixesCount;
this.expectedAverageRideHeight = expectedAverageRideHeight;
this.expectedFixesConsolidated = expectedFixesConsolidated;
}
protected abstract InputStream getInputStream();
@@ -0,0 +1,92 @@
package com.sap.sailing.server.trackfiles.test;
import org.junit.Test;
import com.sap.sailing.server.trackfiles.impl.doublefix.DoubleFixProcessor;
import com.sap.sailing.server.trackfiles.impl.doublefix.DoubleVectorFixData;
import com.sap.sailing.server.trackfiles.impl.doublefix.DownsamplerTo1HzProcessor;
import junit.framework.Assert;
public class DownsamplerTo1HzProcessorTest {
@Test
public void testConsolidation() {
MyProcessor delegate = new MyProcessor();
DownsamplerTo1HzProcessor processor = new DownsamplerTo1HzProcessor(4, delegate);
int second = 1 * 1000;
processor.accept(new DoubleVectorFixData(second + 100, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 300, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 550, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 900, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 999, new double[] { 0d, 0d, 0d, 0d }));
second = 10 * 1000;
processor.accept(new DoubleVectorFixData(second + 130, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 200, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 560, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 600, new double[] { 0d, 0d, 0d, 0d }));
second = 30 * 1000;
processor.accept(new DoubleVectorFixData(second + 130, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 200, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 560, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 561, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 562, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 563, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 564, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 600, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 700, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 800, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 900, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 910, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(second + 920, new double[] { 0d, 0d, 0d, 0d }));
second = 50 * 1000;
processor.accept(new DoubleVectorFixData(second + 130, new double[] { 0d, 0d, 0d, 0d }));
Assert.assertEquals("Count consolidated fixes", 3, delegate.countAccepted);
Assert.assertFalse("Not finished yet", delegate.fineshedWasCalled);
processor.finish();
Assert.assertEquals("Count consolidated fixes", 4, delegate.countAccepted);
Assert.assertTrue("Finished", delegate.fineshedWasCalled);
}
@Test
public void testAverageComputation() {
MyProcessor delegate = new MyProcessor();
DownsamplerTo1HzProcessor processor = new DownsamplerTo1HzProcessor(4, delegate);
int second = 1 * 1000;
processor.accept(new DoubleVectorFixData(second + 100, new double[] { 12.40d, 1d, 1d, 0d }));
processor.accept(new DoubleVectorFixData(second + 300, new double[] { 0.01d, 2d, 1d, 0d }));
processor.accept(new DoubleVectorFixData(second + 550, new double[] { -10.00d, 3d, 1d, 0d }));
processor.accept(new DoubleVectorFixData(second + 900, new double[] { 1.00d, 2d, 1d, 0d }));
processor.accept(new DoubleVectorFixData(second + 999, new double[] { 0.00d, 1d, 1d, 0d }));
processor.finish();
Assert.assertTrue("Finished", delegate.fineshedWasCalled);
Assert.assertEquals("Count consolidated fixes", 1, delegate.countAccepted);
DoubleVectorFixData lastFix = delegate.lastFix;
Assert.assertNotNull("Lastfix", lastFix);
Assert.assertEquals("Avg col 1", (12.4d + 0.01d - 10d + 1d + 0d) / 5d, lastFix.getFix()[0]);
Assert.assertEquals("Avg col 2", 9d / 5d, lastFix.getFix()[1]);
Assert.assertEquals("Avg col 3", 1d, lastFix.getFix()[2]);
Assert.assertEquals("Avg col 4", 0d, lastFix.getFix()[3]);
}
private final class MyProcessor implements DoubleFixProcessor {
public int countAccepted = 0;
public boolean fineshedWasCalled = false;
public DoubleVectorFixData lastFix;
@Override
public void finish() {
fineshedWasCalled = true;
}
@Override
public void accept(DoubleVectorFixData fix) {
lastFix = fix;
countAccepted++;
}
}
}
@@ -0,0 +1,116 @@
package com.sap.sailing.server.trackfiles.test;
import java.util.UUID;
import org.junit.Test;
import com.sap.sailing.domain.common.tracking.DoubleVectorFix;
import com.sap.sailing.domain.trackfiles.TrackFileImportDeviceIdentifier;
import com.sap.sailing.domain.trackimport.DoubleVectorFixImporter;
import com.sap.sailing.server.trackfiles.impl.doublefix.DoubleVectorFixData;
import com.sap.sailing.server.trackfiles.impl.doublefix.LearningBatchProcessor;
import com.sap.sse.common.TimePoint;
import com.sap.sse.common.Util;
import com.sap.sse.common.impl.MillisecondsTimePoint;
import junit.framework.Assert;
public class LearningBatchProcessorTest {
@Test
public void testConsolidation() {
MyCallback delegate = new MyCallback();
LearningBatchProcessor processor = new LearningBatchProcessor(5, 10, delegate, device);
int currentSecondAsMillis = 1 * 1000;
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 100, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 300, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 550, new double[] { 0d, 0d, 0d, 0d }));
Assert.assertEquals("Batch not full", 0, delegate.batchesDelivered);
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 900, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 999, new double[] { 0d, 0d, 0d, 0d }));
Assert.assertEquals("Batch not full, should be learing", 0, delegate.batchesDelivered);
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 130, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 200, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 560, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 600, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 130, new double[] { 0d, 0d, 0d, 0d }));
Assert.assertEquals("Batches after learning", 2, delegate.batchesDelivered);
Assert.assertEquals("Last batch after learning", 5, delegate.lastBatchSize);
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 200, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 560, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 561, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 562, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 563, new double[] { 0d, 0d, 0d, 0d }));
Assert.assertEquals("Next batch received", 3, delegate.batchesDelivered);
Assert.assertEquals("First batch after learning", 5, delegate.lastBatchSize);
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 564, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 600, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 700, new double[] { 0d, 0d, 0d, 0d }));
processor.accept(new DoubleVectorFixData(currentSecondAsMillis + 800, new double[] { 0d, 0d, 0d, 0d }));
Assert.assertEquals("Still on same batch", 3, delegate.batchesDelivered);
processor.finish();
Assert.assertEquals("Batch count after finish", 4, delegate.batchesDelivered);
Assert.assertEquals("Last batch size after finish", 4, delegate.lastBatchSize);
Assert.assertEquals("Last fix ", (currentSecondAsMillis + 800), delegate.lastFix.getTimePoint().asMillis());
}
private final class MyCallback implements DoubleVectorFixImporter.Callback {
public int lastBatchSize = 0;
public int batchesDelivered = 0;
public DoubleVectorFix lastFix;
@Override
public void addFixes(Iterable<DoubleVectorFix> fix, TrackFileImportDeviceIdentifier device) {
lastBatchSize = Util.size(fix);
batchesDelivered++;
lastFix = Util.last(fix);
}
}
TrackFileImportDeviceIdentifier device = new TrackFileImportDeviceIdentifier() {
private static final long serialVersionUID = 870061505508090216L;
private UUID uuid = UUID.randomUUID();
private MillisecondsTimePoint uploadedAt = new MillisecondsTimePoint(0);
@Override
public String getStringRepresentation() {
return "stringRepresentation";
}
@Override
public String getIdentifierType() {
return "identifierType";
}
@Override
public TimePoint getUploadedAt() {
return uploadedAt;
}
@Override
public String getTrackName() {
return "trackname";
}
@Override
public UUID getId() {
return uuid;
}
@Override
public String getFileName() {
return "filename";
}
};
}