diff --git a/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/sensordata/BravoSensorDataMetadata.java b/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/sensordata/BravoSensorDataMetadata.java index 4800af20318..b7b1b9dd201 100644 --- a/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/sensordata/BravoSensorDataMetadata.java +++ b/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/sensordata/BravoSensorDataMetadata.java @@ -5,49 +5,79 @@ import java.util.Collections; import java.util.List; /** - * Metadata that defines the column structure of {@link com.sap.sailing.domain.common.tracking.DoubleVectorFix}es when imported. + * Metadata that defines the column structure of {@link com.sap.sailing.domain.common.tracking.DoubleVectorFix}es when + * imported. + * + * The current implementation only stores a subset of the information available during the import. */ public enum BravoSensorDataMetadata { + INSTANCE; - private static final String RIDE_HEIGHT = "RideHeight"; private static final String RIDE_HEIGHT_PORT_HULL = "RideHeightPortHull"; private static final String RIDE_HEIGHT_STBD_HULL = "RideHeightStbdHull"; - private static final String YAW = "ImuSensor_Yaw"; - private static final String ROLL = "ImuSensor_Roll"; + private static final String HEEL = "Heel"; private static final String PITCH = "ImuSensor_Pitch"; + private final int HEADER_COLUMN_OFFSET = 3; - private final List columns = Collections.unmodifiableList(Arrays.asList(RIDE_HEIGHT, RIDE_HEIGHT_PORT_HULL, - RIDE_HEIGHT_STBD_HULL, "Heel", "Trim", "ImuSensor_GyroX", "ImuSensor_GyroY", "ImuSensor_GyroZ", - PITCH, ROLL, YAW, "ImuSensor_LinearAccX", "ImuSensor_LinearAccY", - "ImuSensor_LinearAccZ", "Hb_Z", "Dn_Z", "Db_Z", "LKF_ride_hgh", "LKF_ride_hgh_Position", - "LKF_ride_hgh_Velocity", "LKF_ride_hgh_Acceleration", "LKF_ride_hgh_PositionError", - "LKF_ride_hgh_VelocityError", "LKF_ride_hgh_AccelerationError", "Gps_Gga_PosFixTime", "Gps_Gga_Lat", - "Gps_Gga_Lon", "Gps_Gga_QI", "Gps_Gga_HDOP", "Gps_Gga_AntHeight", "Gps_Vtg_TMG", "Gps_Vtg_SOGKnots", - "Gps_Rmc_MagVar", "Gps_EastVelocity", "Gps_NorthVelocity", "Gps_UpVelocity", "BravoNet_Node0x0A0_Ch0", - "BravoNet_Node0x0A0_Voltage", "BravoNet_Node0x0A0_Temperature", "BravoNet_Node0x0A0_Current")); - - public final int rideHeightColumn = getColumnIndex(RIDE_HEIGHT); - public final int rideHeightPortHullColumn = getColumnIndex(RIDE_HEIGHT_PORT_HULL); - public final int rideHeightStarboardHullColumn = getColumnIndex(RIDE_HEIGHT_STBD_HULL); - public final int pitchColumn = getColumnIndex(PITCH); - public final int rollColumn = getColumnIndex(ROLL); - public final int yawColumn = getColumnIndex(YAW); - public final int columnCount = columns.size(); + /** + * + * Available information in file:
+ * The columns field defines which information will be loaded into the system + * + */ + private final List fileColumns = Collections.unmodifiableList(Arrays.asList( // + "RideHeight", RIDE_HEIGHT_PORT_HULL, RIDE_HEIGHT_STBD_HULL, HEEL, "Trim", "ImuSensor_GyroX", + "ImuSensor_GyroY", "ImuSensor_GyroZ", PITCH, "ImuSensor_Roll", "ImuSensor_Yaw", "ImuSensor_LinearAccX", + "ImuSensor_LinearAccY", "ImuSensor_LinearAccZ", + "Hb_Z", "Dn_Z", "Db_Z", "LKF_ride_hgh", "LKF_ride_hgh_Position", "LKF_ride_hgh_Velocity", + "LKF_ride_hgh_Acceleration", "LKF_ride_hgh_PositionError", "LKF_ride_hgh_VelocityError", + "LKF_ride_hgh_AccelerationError", "Gps_Gga_PosFixTime", "Gps_Gga_Lat", "Gps_Gga_Lon", "Gps_Gga_QI", + "Gps_Gga_HDOP", "Gps_Gga_AntHeight", "Gps_Vtg_TMG", "Gps_Vtg_SOGKnots", "Gps_Rmc_MagVar", "Gps_EastVelocity", + "Gps_NorthVelocity", "Gps_UpVelocity", "BravoNet_Node0x0A0_Ch0", "BravoNet_Node0x0A0_Voltage", + "BravoNet_Node0x0A0_Temperature", "BravoNet_Node0x0A0_Current" + )); - public List getColumns() { - return columns; + private final List trackColumns = Collections.unmodifiableList(Arrays.asList( // + RIDE_HEIGHT_PORT_HULL, // + RIDE_HEIGHT_STBD_HULL, // + HEEL, // + PITCH // + )); + + public final int trackColumnCount = trackColumns.size(); + public final int rideHeightPortHullColumn = getTrackColumnIndex(RIDE_HEIGHT_PORT_HULL); + public final int rideHeightStarboardHullColumn = getTrackColumnIndex(RIDE_HEIGHT_STBD_HULL); + public final int heelColumn = getTrackColumnIndex(HEEL); + public final int pitchColumn = getTrackColumnIndex(PITCH); + + /** + * The fields available in the import file + * + * @return + */ + public List getFileColumns() { + return fileColumns; } - public boolean hasColumn(String columnName) { - return columns.contains(columnName); + /** + * The fields to be loaded into a track + * + * @return + */ + public List getTrackColumns() { + return trackColumns; } - public int getColumnIndex(String columnName) { - return columns.indexOf(columnName); + public boolean trackHasColumn(String columnName) { + return trackColumns.contains(columnName); } - + + public int getTrackColumnIndex(String columnName) { + return trackColumns.indexOf(columnName); + } + public int getHeaderColumnOffset() { return HEADER_COLUMN_OFFSET; } diff --git a/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/tracking/BravoFix.java b/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/tracking/BravoFix.java index 3b441386140..2f5e69b6724 100644 --- a/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/tracking/BravoFix.java +++ b/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/tracking/BravoFix.java @@ -33,9 +33,7 @@ public interface BravoFix extends SensorFix { @Statistic(messageKey="pitch", resultDecimals=1) double getPitch(); - @Statistic(messageKey="yaw", resultDecimals=1) - double getYaw(); + @Statistic(messageKey = "heel", resultDecimals = 1) + double getHeel(); - @Statistic(messageKey="roll", resultDecimals=1) - double getRoll(); } diff --git a/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/tracking/impl/BravoFixImpl.java b/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/tracking/impl/BravoFixImpl.java index a70b783746a..7a612a9551c 100644 --- a/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/tracking/impl/BravoFixImpl.java +++ b/java/com.sap.sailing.domain.common/src/com/sap/sailing/domain/common/tracking/impl/BravoFixImpl.java @@ -6,23 +6,39 @@ import com.sap.sailing.domain.common.sensordata.BravoSensorDataMetadata; import com.sap.sailing.domain.common.tracking.BravoFix; import com.sap.sailing.domain.common.tracking.DoubleVectorFix; import com.sap.sse.common.TimePoint; -import com.sap.sse.common.Util; /** * Implementation of {@link BravoFix} that wraps a {@link DoubleVectorFix} which holds the actual sensor data. */ public class BravoFixImpl implements BravoFix { private static final long serialVersionUID = 2033254212013221160L; - private final DoubleVectorFix fix; + private final MeterDistance computedRideHeight; + private final MeterDistance rideHeightPortHull; + private final MeterDistance rideHeightStarboardHull; + private final double pitch; + private final double heel; + private final boolean computedIsFoiling; public BravoFixImpl(DoubleVectorFix fix) { this.fix = fix; + + pitch = fix.get(BravoSensorDataMetadata.INSTANCE.pitchColumn); + heel = fix.get(BravoSensorDataMetadata.INSTANCE.heelColumn); + double rideHeightPortHullasDouble = fix.get(BravoSensorDataMetadata.INSTANCE.rideHeightPortHullColumn); + double rideHeightStarboardHullasDouble = fix + .get(BravoSensorDataMetadata.INSTANCE.rideHeightStarboardHullColumn); + + rideHeightPortHull = new MeterDistance(rideHeightPortHullasDouble); + rideHeightStarboardHull = new MeterDistance(rideHeightStarboardHullasDouble); + + computedRideHeight = new MeterDistance(Math.min(rideHeightPortHullasDouble, rideHeightStarboardHullasDouble)); + computedIsFoiling = computedRideHeight.compareTo(MIN_FOILING_HEIGHT_THRESHOLD) >= 0; } @Override public double get(String valueName) { - int index = BravoSensorDataMetadata.INSTANCE.getColumnIndex(valueName); + int index = BravoSensorDataMetadata.INSTANCE.getTrackColumnIndex(valueName); if (index < 0) { throw new IllegalArgumentException("Unknown value \"" + valueName + "\" for " + getClass().getSimpleName()); } @@ -36,36 +52,31 @@ public class BravoFixImpl implements BravoFix { @Override public Distance getRideHeight() { - return Util.min(getRideHeightPortHull(), getRideHeightStarboardHull()); + return computedRideHeight; } - + @Override public Distance getRideHeightPortHull() { - return new MeterDistance(fix.get(BravoSensorDataMetadata.INSTANCE.rideHeightPortHullColumn)); + return rideHeightPortHull; } - + @Override public Distance getRideHeightStarboardHull() { - return new MeterDistance(fix.get(BravoSensorDataMetadata.INSTANCE.rideHeightStarboardHullColumn)); + return rideHeightStarboardHull; } - + @Override public boolean isFoiling() { - return getRideHeight().compareTo(MIN_FOILING_HEIGHT_THRESHOLD) >= 0; + return computedIsFoiling; } @Override public double getPitch() { - return fix.get(BravoSensorDataMetadata.INSTANCE.pitchColumn); + return pitch; } - + @Override - public double getYaw() { - return fix.get(BravoSensorDataMetadata.INSTANCE.yawColumn); - } - - @Override - public double getRoll() { - return fix.get(BravoSensorDataMetadata.INSTANCE.rollColumn); + public double getHeel() { + return heel; } } diff --git a/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/SensorFixStoreAndLoadTest.java b/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/SensorFixStoreAndLoadTest.java index 5ce6b2471ba..2e7f4e17d99 100644 --- a/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/SensorFixStoreAndLoadTest.java +++ b/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/SensorFixStoreAndLoadTest.java @@ -513,8 +513,7 @@ public class SensorFixStoreAndLoadTest { } private DoubleVectorFix createBravoDoubleVectorFixWithRideHeight(long timestamp, double rideHeight) { - double[] fixData = new double[BravoSensorDataMetadata.INSTANCE.columnCount]; - fixData[BravoSensorDataMetadata.INSTANCE.rideHeightColumn] = rideHeight; + double[] fixData = new double[BravoSensorDataMetadata.INSTANCE.trackColumnCount]; // fill the port/starboard columns as well because their minimum defines the true ride height fixData[BravoSensorDataMetadata.INSTANCE.rideHeightPortHullColumn] = rideHeight; fixData[BravoSensorDataMetadata.INSTANCE.rideHeightStarboardHullColumn] = rideHeight; diff --git a/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/SensorFixStoreTest.java b/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/SensorFixStoreTest.java index 0154d127bfb..71a87ca8638 100644 --- a/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/SensorFixStoreTest.java +++ b/java/com.sap.sailing.domain.racelogtrackingadapter.test/src/com/sap/sailing/domain/racelogtracking/test/impl/SensorFixStoreTest.java @@ -330,8 +330,9 @@ public class SensorFixStoreTest { } private DoubleVectorFix createBravoDoubleVectorFixWithRideHeight(long timestamp, double rideHeight) { - double[] fixData = new double[BravoSensorDataMetadata.INSTANCE.columnCount]; - fixData[BravoSensorDataMetadata.INSTANCE.rideHeightColumn] = rideHeight; + double[] fixData = new double[BravoSensorDataMetadata.INSTANCE.trackColumnCount]; + fixData[BravoSensorDataMetadata.INSTANCE.rideHeightPortHullColumn] = rideHeight; + fixData[BravoSensorDataMetadata.INSTANCE.rideHeightStarboardHullColumn] = rideHeight; return new DoubleVectorFixImpl(new MillisecondsTimePoint(timestamp), fixData); } } diff --git a/java/com.sap.sailing.domain.racelogtrackingadapter/src/com/sap/sailing/domain/racelogtracking/impl/BravoDataFixMapper.java b/java/com.sap.sailing.domain.racelogtrackingadapter/src/com/sap/sailing/domain/racelogtracking/impl/BravoDataFixMapper.java index 45b8e031dac..4469caa01f8 100644 --- a/java/com.sap.sailing.domain.racelogtrackingadapter/src/com/sap/sailing/domain/racelogtracking/impl/BravoDataFixMapper.java +++ b/java/com.sap.sailing.domain.racelogtrackingadapter/src/com/sap/sailing/domain/racelogtracking/impl/BravoDataFixMapper.java @@ -32,4 +32,4 @@ public class BravoDataFixMapper implements SensorFixMapper> eventType) { return RegattaLogDeviceCompetitorBravoMappingEventImpl.class.isAssignableFrom(eventType); } -} +} \ No newline at end of file diff --git a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/tracking/impl/BravoFixTrackSerializationTest.java b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/tracking/impl/BravoFixTrackSerializationTest.java index 4b855831052..13d991266a5 100644 --- a/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/tracking/impl/BravoFixTrackSerializationTest.java +++ b/java/com.sap.sailing.domain.test/src/com/sap/sailing/domain/tracking/impl/BravoFixTrackSerializationTest.java @@ -116,8 +116,7 @@ public class BravoFixTrackSerializationTest { } private void addOrReplaceBravoFixToTrack(boolean replace) { - double[] fixData = new double[BravoSensorDataMetadata.INSTANCE.columnCount]; - fixData[BravoSensorDataMetadata.INSTANCE.rideHeightColumn] = rideHeight; + double[] fixData = new double[BravoSensorDataMetadata.INSTANCE.trackColumnCount]; // fill the port/starboard columns as well because their minimum defines the true ride height fixData[BravoSensorDataMetadata.INSTANCE.rideHeightPortHullColumn] = rideHeight; fixData[BravoSensorDataMetadata.INSTANCE.rideHeightStarboardHullColumn] = rideHeight; diff --git a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/BravoFixTrackImpl.java b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/BravoFixTrackImpl.java index 790fc324bef..6e90953b125 100644 --- a/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/BravoFixTrackImpl.java +++ b/java/com.sap.sailing.domain/src/com/sap/sailing/domain/tracking/impl/BravoFixTrackImpl.java @@ -28,7 +28,7 @@ public class BravoFixTrackImpl extends S * the name of the track by which it can be obtained from the {@link TrackedRace}. */ public BravoFixTrackImpl(ItemType trackedItem, String trackName) { - super(trackedItem, trackName, BravoSensorDataMetadata.INSTANCE.getColumns(), + super(trackedItem, trackName, BravoSensorDataMetadata.INSTANCE.getFileColumns(), BravoFixTrack.TRACK_NAME + " for " + trackedItem); } diff --git a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/MasterDataImportTest.java b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/MasterDataImportTest.java index 0bc2f6ae15a..dd14c50d80e 100644 --- a/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/MasterDataImportTest.java +++ b/java/com.sap.sailing.server.test/src/com/sap/sailing/server/test/MasterDataImportTest.java @@ -307,9 +307,10 @@ public class MasterDataImportTest { logTimePoint4, author, competitor, deviceIdentifier2, logTimePoint, logTimePoint3); regatta.getRegattaLog().add(bravoMappingEvent); - double[] fixData = new double[BravoSensorDataMetadata.INSTANCE.columnCount]; + double[] fixData = new double[BravoSensorDataMetadata.INSTANCE.trackColumnCount]; double rideHeightValue = 1337.0; - fixData[BravoSensorDataMetadata.INSTANCE.rideHeightColumn] = rideHeightValue; + fixData[BravoSensorDataMetadata.INSTANCE.rideHeightPortHullColumn] = rideHeightValue; + fixData[BravoSensorDataMetadata.INSTANCE.rideHeightStarboardHullColumn] = rideHeightValue; DoubleVectorFixImpl doubleVectorFix = new DoubleVectorFixImpl(logTimePoint2, fixData); sourceService.getSensorFixStore().storeFix(deviceIdentifier2, doubleVectorFix); diff --git a/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/BravoDataImportTest.java b/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/BravoDataImportTest.java index d0083b0c9de..064f35a4b2d 100644 --- a/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/BravoDataImportTest.java +++ b/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/BravoDataImportTest.java @@ -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(); diff --git a/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/DownsamplerTo1HzProcessorTest.java b/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/DownsamplerTo1HzProcessorTest.java new file mode 100644 index 00000000000..80f1f334bec --- /dev/null +++ b/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/DownsamplerTo1HzProcessorTest.java @@ -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++; + } + } + +} diff --git a/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/LearningBatchProcessorTest.java b/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/LearningBatchProcessorTest.java new file mode 100644 index 00000000000..cf3127859ba --- /dev/null +++ b/java/com.sap.sailing.server.trackfiles.test/src/com/sap/sailing/server/trackfiles/test/LearningBatchProcessorTest.java @@ -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 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"; + } + }; +} diff --git a/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/BravoDataImporterImpl.java b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/BravoDataImporterImpl.java index 6ca392f137d..f18665cde9d 100644 --- a/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/BravoDataImporterImpl.java +++ b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/BravoDataImporterImpl.java @@ -5,12 +5,6 @@ import java.io.IOException; import java.io.InputStream; import java.io.InputStreamReader; import java.io.Serializable; -import java.time.Instant; -import java.time.LocalDateTime; -import java.time.ZoneOffset; -import java.time.format.DateTimeFormatter; -import java.util.ArrayList; -import java.util.Date; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -18,7 +12,6 @@ import java.util.UUID; import java.util.concurrent.atomic.AtomicLong; import java.util.logging.Level; import java.util.logging.Logger; -import java.util.stream.IntStream; import java.util.zip.GZIPInputStream; import com.sap.sailing.domain.abstractlog.AbstractLogEventAuthor; @@ -26,13 +19,15 @@ import com.sap.sailing.domain.abstractlog.regatta.events.RegattaLogDeviceCompeti import com.sap.sailing.domain.abstractlog.regatta.events.impl.RegattaLogDeviceCompetitorBravoMappingEventImpl; import com.sap.sailing.domain.base.Competitor; import com.sap.sailing.domain.common.sensordata.BravoSensorDataMetadata; -import com.sap.sailing.domain.common.tracking.DoubleVectorFix; -import com.sap.sailing.domain.common.tracking.impl.DoubleVectorFixImpl; import com.sap.sailing.domain.racelogtracking.DeviceIdentifier; import com.sap.sailing.domain.trackfiles.TrackFileImportDeviceIdentifier; import com.sap.sailing.domain.trackfiles.TrackFileImportDeviceIdentifierImpl; import com.sap.sailing.domain.trackimport.DoubleVectorFixImporter; import com.sap.sailing.domain.trackimport.FormatNotSupportedException; +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 com.sap.sailing.server.trackfiles.impl.doublefix.LearningBatchProcessor; import com.sap.sse.common.TimePoint; import com.sap.sse.common.impl.MillisecondsTimePoint; @@ -42,10 +37,8 @@ import com.sap.sse.common.impl.MillisecondsTimePoint; */ public class BravoDataImporterImpl implements DoubleVectorFixImporter { private final Logger LOG = Logger.getLogger(DoubleVectorFixImporter.class.getName()); - private final static DateTimeFormatter DATE_FMT = DateTimeFormatter.ofPattern("yyyyMMdd.HHmmss.SSSSSS"); private final BravoSensorDataMetadata metadata = BravoSensorDataMetadata.INSTANCE; private final String BOF = "jjlDATE\tjjlTIME\tEpoch"; - private static final int BATCH_SIZE = 5000; public void importFixes(InputStream inputStream, Callback callback, final String filename, String sourceName) throws FormatNotSupportedException, IOException { @@ -79,22 +72,15 @@ public class BravoDataImporterImpl implements DoubleVectorFixImporter { } LOG.fine("Validate and parse header columns"); final Map colIndices = validateAndParseHeader(headerLine); - LOG.fine("Parse and store data in batches of " + BATCH_SIZE + " items"); - final ArrayList collectedFixes = new ArrayList<>(BATCH_SIZE); + + DoubleFixProcessor downsampler = createProcessor(metadata, callback, trackIdentifier); + buffer.lines().forEach(line -> { lineNr.incrementAndGet(); - DoubleVectorFixImpl vectorFix = parseLine(lineNr.get(), filename, line, colIndices); - if (vectorFix != null) { - collectedFixes.add(vectorFix); - if (collectedFixes.size() == BATCH_SIZE) { - callback.addFixes(collectedFixes, trackIdentifier); - collectedFixes.clear(); - } - } + downsampler.accept(parseLine(lineNr.get(), filename, line, colIndices)); }); - if (!collectedFixes.isEmpty()) { - callback.addFixes(collectedFixes, trackIdentifier); - } + + downsampler.finish(); buffer.close(); } } catch (Exception e) { @@ -102,40 +88,46 @@ public class BravoDataImporterImpl implements DoubleVectorFixImporter { } } + /** + * This method creates the double fix processor chain used to downsample and batch the stream of fixes parsed by the + * importer. + * + * This method is protected so it can be overridden by the test case. + * + * @param metadata + * @param callback + * @param trackIdentifier + * @return + */ + protected DoubleFixProcessor createProcessor(BravoSensorDataMetadata metadata, Callback callback, + final TrackFileImportDeviceIdentifier trackIdentifier) { + LearningBatchProcessor batchProcessor = new LearningBatchProcessor(5000, 5000, callback, trackIdentifier); + DoubleFixProcessor downsampler = new DownsamplerTo1HzProcessor(metadata.trackColumnCount, batchProcessor); + return downsampler; + } + /** * Parses the CSV line and reads the double data values in the order defined by the col enums. */ - private DoubleVectorFixImpl parseLine(long lineNr, String filename, String line, Map colIndices) { + private DoubleVectorFixData parseLine(long lineNr, String filename, String line, Map columnsInFile) { try { - String[] contentTokens = split(line); - TimePoint fixTp; - String epochColValue = contentTokens[2]; + String[] fileContentTokens = split(line); + String epochColValue = fileContentTokens[2]; + long epoch; if (epochColValue != null && epochColValue.length() > 0) { epochColValue = epochColValue.substring(0, epochColValue.indexOf(".")); - long epoch = Long.parseLong(epochColValue); - fixTp = new MillisecondsTimePoint(epoch); + epoch = Long.parseLong(epochColValue); } else { - String jjLDATE = contentTokens[0].substring(0, contentTokens[0].indexOf(".")); - StringBuilder dtb = new StringBuilder(jjLDATE); - String jjlTIME = contentTokens[1]; - int offset = 6 - jjlTIME.indexOf("."); - IntStream.range(0, offset).forEach(n -> dtb.append("0")); - dtb.append(jjlTIME); - LocalDateTime day = DATE_FMT.parse(dtb.toString(), LocalDateTime::from); - Instant instant = day.toInstant(ZoneOffset.UTC); - fixTp = new MillisecondsTimePoint(Date.from(instant)); + // we don't have epoch, skip the line + return null; } - // code for time adjustment for foiling test data - // long r16TP = ZonedDateTime.of(LocalDateTime.of(2016, 3, 19, 11, 55), ZoneId.of("Europe/Berlin")) - // .toEpochSecond() * 1000; - // fixTp = fixTp.plus(r16TP - 1461828459589l); - double[] fixData = new double[metadata.columnCount]; - for (int columnIndexInFix = 0; columnIndexInFix < fixData.length; columnIndexInFix++) { - String columnName = metadata.getColumns().get(columnIndexInFix); - Integer columnIndexInFile = colIndices.get(columnName); - fixData[columnIndexInFix] = Double.parseDouble(contentTokens[columnIndexInFile]); + double[] trackFixData = new double[metadata.trackColumnCount]; + for (int trackColumnIdx = 0; trackColumnIdx < metadata.trackColumnCount; trackColumnIdx++) { + String columnNameToSearchForInFile = metadata.getTrackColumns().get(trackColumnIdx); + Integer columnsInFileIdx = columnsInFile.get(columnNameToSearchForInFile); + trackFixData[trackColumnIdx] = Double.parseDouble(fileContentTokens[columnsInFileIdx]); } - return new DoubleVectorFixImpl(fixTp, fixData); + return new DoubleVectorFixData(epoch, trackFixData); } catch (Exception e) { LOG.warning( "Error parsing line nr " + lineNr + " in file " + filename + "with exception: " + e.getMessage()); @@ -150,7 +142,7 @@ public class BravoDataImporterImpl implements DoubleVectorFixImporter { String header = headerTokens[j]; colIndicesInFile.put(header, j); } - List columnsInFix = metadata.getColumns(); + List columnsInFix = metadata.getFileColumns(); if (colIndicesInFile.size() != columnsInFix.size() || !colIndicesInFile.keySet().containsAll(columnsInFix)) { LOG.log(Level.SEVERE, "Missing headers"); throw new RuntimeException("Missing headers in import files"); diff --git a/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/DoubleFixProcessor.java b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/DoubleFixProcessor.java new file mode 100644 index 00000000000..661b298b8f1 --- /dev/null +++ b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/DoubleFixProcessor.java @@ -0,0 +1,7 @@ +package com.sap.sailing.server.trackfiles.impl.doublefix; + + +public interface DoubleFixProcessor { + void accept(DoubleVectorFixData fix); + void finish(); +} \ No newline at end of file diff --git a/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/DoubleVectorFixData.java b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/DoubleVectorFixData.java new file mode 100644 index 00000000000..9c30e56b839 --- /dev/null +++ b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/DoubleVectorFixData.java @@ -0,0 +1,28 @@ +package com.sap.sailing.server.trackfiles.impl.doublefix; + +public final class DoubleVectorFixData { + private long timepointInMs; + private double[] fix; + + public DoubleVectorFixData(long timepoint, double[] fix) { + super(); + this.timepointInMs = timepoint; + this.fix = fix; + } + + public void correctTimepointBy(long offset) { + timepointInMs += offset; + } + + public long getTimepointInMs() { + return timepointInMs; + } + + public long getFixSecond() { + return timepointInMs / 1000; + } + + public double[] getFix() { + return fix; + } +} \ No newline at end of file diff --git a/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/DownsamplerTo1HzProcessor.java b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/DownsamplerTo1HzProcessor.java new file mode 100644 index 00000000000..ecdd5955f85 --- /dev/null +++ b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/DownsamplerTo1HzProcessor.java @@ -0,0 +1,78 @@ +package com.sap.sailing.server.trackfiles.impl.doublefix; + +import java.util.ArrayList; +import java.util.logging.Logger; + +import com.sap.sailing.domain.trackimport.DoubleVectorFixImporter; + +/** + * This class consolidates fixes in a sub-second range to one fix per second and computes the average of all values. + * + */ +public final class DownsamplerTo1HzProcessor implements DoubleFixProcessor { + private final Logger LOG = Logger.getLogger(DoubleVectorFixImporter.class.getName()); + final private ArrayList fixesInTheCurrentSecond = new ArrayList<>(); + final private int nrOfColumsInTrack; + final private DoubleFixProcessor delegateProcesssor; + private long currentSecond = 0; + private long currentOffsetInMs = 0; + private long countSourceTtl = 0; + private long countImportedTtl = 0; + + public DownsamplerTo1HzProcessor(int nrOfColumsInTrack, + DoubleFixProcessor consumerOfConsolidatedFixes) { + this.nrOfColumsInTrack = nrOfColumsInTrack; + this.delegateProcesssor = consumerOfConsolidatedFixes; + } + + @Override + public void accept(DoubleVectorFixData fix) { + if (fix == null) + return; + fix.correctTimepointBy(currentOffsetInMs); + if (fix.getFixSecond() < currentSecond) { + currentOffsetInMs = (currentSecond - fix.getFixSecond() + 1) * 1000; + LOG.warning("Timepoint before last second, using offser of " + currentOffsetInMs + "ms from now on"); + } + if (currentSecond != fix.getFixSecond()) { + computeDownsampledFixForCurrentSecond(); + currentSecond = fix.getFixSecond(); + } + countSourceTtl++; + fixesInTheCurrentSecond.add(fix); + } + + public void finish() { + computeDownsampledFixForCurrentSecond(); + delegateProcesssor.finish(); + LOG.fine("Imported " + countImportedTtl + " fixes from " + countSourceTtl); + } + + private void computeDownsampledFixForCurrentSecond() { + final double[] computedAverage = new double[nrOfColumsInTrack]; + final int numberOfFixesInSecond = fixesInTheCurrentSecond.size(); + if (numberOfFixesInSecond == 0) { + return; + } + for (DoubleVectorFixData d : fixesInTheCurrentSecond) { + final double[] fix = d.getFix(); + for (int colIdx = 0; colIdx < nrOfColumsInTrack; colIdx++) { + computedAverage[colIdx] += fix[colIdx]; + } + } + fixesInTheCurrentSecond.clear(); + for (int colIdx = 0; colIdx < nrOfColumsInTrack; colIdx++) { + computedAverage[colIdx] /= (double) numberOfFixesInSecond; + } + delegateProcesssor.accept(new DoubleVectorFixData(currentSecond * 1000 + 500, computedAverage)); + countImportedTtl++; + } + + public long getCountImportedTtl() { + return countImportedTtl; + } + + public long getCountSourceTtl() { + return countSourceTtl; + } +} \ No newline at end of file diff --git a/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/LearningBatchProcessor.java b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/LearningBatchProcessor.java new file mode 100644 index 00000000000..5ccb98a982f --- /dev/null +++ b/java/com.sap.sailing.server.trackfiles/src/com/sap/sailing/server/trackfiles/impl/doublefix/LearningBatchProcessor.java @@ -0,0 +1,90 @@ +package com.sap.sailing.server.trackfiles.impl.doublefix; + +import java.util.ArrayList; + +import com.sap.sailing.domain.common.tracking.DoubleVectorFix; +import com.sap.sailing.domain.common.tracking.impl.DoubleVectorFixImpl; +import com.sap.sailing.domain.trackfiles.TrackFileImportDeviceIdentifier; +import com.sap.sailing.domain.trackimport.DoubleVectorFixImporter.Callback; +import com.sap.sse.common.impl.MillisecondsTimePoint; + +public final class LearningBatchProcessor implements DoubleFixProcessor { + private final int batchSize; + private final ArrayList learnedFixes; + private final ArrayList collectedFixes; + private final Callback callback; + private final TrackFileImportDeviceIdentifier deviceIdentifier; + private boolean isLearning = true; + private long fixesToLearn; + + public LearningBatchProcessor(int batchSize, int fixesToLearn, Callback callback, + TrackFileImportDeviceIdentifier deviceIdentifier) { + super(); + this.batchSize = batchSize; + this.fixesToLearn = fixesToLearn; + this.collectedFixes = new ArrayList<>(batchSize); + this.learnedFixes = new ArrayList<>(fixesToLearn); + this.callback = callback; + this.deviceIdentifier = deviceIdentifier; + } + + @Override + public void accept(DoubleVectorFixData fix) { + if (isLearning) { + learnFix(fix); + } else { + processFix(fix); + } + } + + /** + * In the learn phase the fixes just get collected for learning. + * + * @param fix + */ + private void learnFix(DoubleVectorFixData fix) { + learnedFixes.add(fix); + + if (learnedFixes.size() >= fixesToLearn) { + finishLearning(); + } + } + + private void finishLearning() { + // Now we have enough data to process + isLearning = false; + // process the loaded learning fixes + for (DoubleVectorFixData doubleVectorFix : learnedFixes) { + processFix(doubleVectorFix); + } + } + + /** + * When processing a fix, data can be corrected before added to the collection of Fixes that will get stored in the + * mongo.... + * + * @param fix + */ + private void processFix(DoubleVectorFixData fix) { + // TODO process fix + collectedFixes.add(new DoubleVectorFixImpl(new MillisecondsTimePoint(fix.getTimepointInMs()), fix.getFix())); + int currentlyCollectedFixes = collectedFixes.size(); + if (currentlyCollectedFixes >= batchSize) { + pushCollectedFixes(); + } + + } + + @Override + public void finish() { + if (isLearning) { + finishLearning(); + } + pushCollectedFixes(); + } + + private void pushCollectedFixes() { + callback.addFixes(collectedFixes, deviceIdentifier); + collectedFixes.clear(); + } +} \ No newline at end of file