bug6236: (WIP) preparing for caching maximum time ranges with same device mappings

This commit is contained in:
Axel Uhl
2026-04-14 17:37:32 +02:00
parent 56d8c77474
commit abb891e787
2 changed files with 208 additions and 150 deletions
@@ -210,158 +210,159 @@ public class FixLoaderAndTracker implements TrackingDataLoader {
Iterable<Timed> fixes, boolean returnManeuverChanges, boolean returnLiveDelay) {
final Set<RegattaAndRaceIdentifier> maneuverChanged = new HashSet<>();
final Map<RegattaAndRaceIdentifier, Duration> delayToLive = new HashMap<>();
// TODO bug6236: how to improve performance of device mappings look-up when we have received multiple fixes from the same device?
if (!preemptiveStopRequested.get() && trackedRace.getStartOfTracking() != null) {
final TimePoint timePoint = fix.getTimePoint();
deviceMappings.forEachMappingOfDeviceIncludingTimePoint(device, fix.getTimePoint(),
new Consumer<DeviceMappingWithRegattaLogEvent<WithID>>() {
@Override
public void accept(DeviceMappingWithRegattaLogEvent<WithID> mapping) {
mapping.getRegattaLogEvent().accept(new MappingEventVisitor() {
@Override
public void visit(RegattaLogDeviceCompetitorSensorDataMappingEvent event) {
recordSensorFixForCompetitor(event.getMappedTo(), event);
}
@Override
public void visit(RegattaLogDeviceBoatSensorDataMappingEvent event) {
final Boat boat = event.getMappedTo();
final Competitor competitor = trackedRace.getCompetitorOfBoat(boat);
if (competitor != null) {
recordSensorFixForCompetitor(competitor, event);
} else {
logger.log(Level.FINE, ()->"Could not record fix for boat because no competitor could be determined. Boat: " + boat);
for (final Timed fix : fixes) {
final TimePoint timePoint = fix.getTimePoint();
deviceMappings.forEachMappingOfDeviceIncludingTimePoint(device, fix.getTimePoint(),
new Consumer<DeviceMappingWithRegattaLogEvent<WithID>>() {
@Override
public void accept(DeviceMappingWithRegattaLogEvent<WithID> mapping) {
mapping.getRegattaLogEvent().accept(new MappingEventVisitor() {
@Override
public void visit(RegattaLogDeviceCompetitorSensorDataMappingEvent event) {
recordSensorFixForCompetitor(event.getMappedTo(), event);
}
}
private void recordSensorFixForCompetitor(Competitor competitor, RegattaLogDeviceMappingEvent<?> event) {
if (!preemptiveStopRequested.get()) {
@SuppressWarnings("unchecked")
SensorFixMapper<SensorFix, DynamicSensorFixTrack<Competitor, SensorFix>, Competitor> mapper = sensorFixMapperFactory
.createCompetitorMapper((Class<? extends RegattaLogDeviceMappingEvent<?>>) event.getClass());
DynamicSensorFixTrack<Competitor, SensorFix> track = mapper.getTrack(trackedRace, competitor);
if (track != null && trackedRace.isWithinStartAndEndOfTracking(fix.getTimePoint())) {
mapper.addFix(track, (DoubleVectorFix) fix);
@Override
public void visit(RegattaLogDeviceBoatSensorDataMappingEvent event) {
final Boat boat = event.getMappedTo();
final Competitor competitor = trackedRace.getCompetitorOfBoat(boat);
if (competitor != null) {
recordSensorFixForCompetitor(competitor, event);
} else {
logger.log(Level.FINE, ()->"Could not record fix for boat because no competitor could be determined. Boat: " + boat);
}
}
private void recordSensorFixForCompetitor(Competitor competitor, RegattaLogDeviceMappingEvent<?> event) {
if (!preemptiveStopRequested.get()) {
@SuppressWarnings("unchecked")
SensorFixMapper<SensorFix, DynamicSensorFixTrack<Competitor, SensorFix>, Competitor> mapper = sensorFixMapperFactory
.createCompetitorMapper((Class<? extends RegattaLogDeviceMappingEvent<?>>) event.getClass());
DynamicSensorFixTrack<Competitor, SensorFix> track = mapper.getTrack(trackedRace, competitor);
if (track != null && trackedRace.isWithinStartAndEndOfTracking(fix.getTimePoint())) {
mapper.addFix(track, (DoubleVectorFix) fix);
if (returnLiveDelay) {
delayToLive.put(trackedRace.getRaceIdentifier(), new MillisecondsDurationImpl(trackedRace.getDelayToLiveInMillis()));
}
}
}
}
@Override
public void visit(RegattaLogDeviceCompetitorMappingEvent event) {
recordForCompetitor(event.getMappedTo());
}
@Override
public void visit(RegattaLogDeviceBoatMappingEvent event) {
final Boat boat = event.getMappedTo();
final Competitor comp = trackedRace.getCompetitorOfBoat(boat);
if (comp != null) {
recordForCompetitor(comp);
} else {
// this is not necessarily something to warn of; while a boat tracker may continuously track
logger.log(Level.FINE,
()->"Could not record fix for boat because no competitor could be determined. Boat: " + boat);
}
}
private void recordForCompetitor(Competitor comp) {
if (!preemptiveStopRequested.get()) {
if (fix instanceof GPSFixMoving) {
// try to record the fix, and only if it was really to the track,
// check for maneuvers; otherwise, the fix may not have been accepted
// by the race or the track, e.g., because the race's end-of-tracking
// comes before the fix's time point
if (trackedRace.recordFix(comp, (GPSFixMoving) fix)) { // TOOD bug6229: this checks the TrackedRace's tracking interval, but for an MDI we'd also want to intersect with Event/Regatta end date if set
if (returnManeuverChanges) {
RegattaAndRaceIdentifier maneuverChangedAnswer = detectIfManeuverChanged(comp);
if (maneuverChangedAnswer != null) {
maneuverChanged.add(maneuverChangedAnswer);
}
}
if (returnLiveDelay) {
delayToLive.put(trackedRace.getRaceIdentifier(), new MillisecondsDurationImpl(trackedRace.getDelayToLiveInMillis()));
}
}
} else {
logger.log(Level.WARNING,
String.format(
"Could not add fix for competitor (%s) in race (%s), as it"
+ " is no GPSFixMoving, meaning it is missing COG/SOG values",
comp, trackedRace.getRace().getName()));
}
}
}
@Override
public void visit(RegattaLogDeviceMarkMappingEvent event) {
if (!preemptiveStopRequested.get()) {
Mark mark = event.getMappedTo();
final DynamicGPSFixTrack<Mark, GPSFix> markTrack = trackedRace.getOrCreateTrack(mark);
final GPSFix firstFixAtOrAfter;
final boolean forceFix;
if (trackedRace.isWithinStartAndEndOfTracking(fix.getTimePoint())) {
forceFix = false;
} else {
markTrack.lockForRead();
try {
if (Util.isEmpty(markTrack.getRawFixes())
|| (firstFixAtOrAfter = markTrack.getFirstFixAtOrAfter(timePoint)) != null
&& firstFixAtOrAfter.getTimePoint().equals(timePoint)) {
// either the first fix or overwriting an existing one
forceFix = true;
} else {
// checking if the given fix is "better" than an existing one
TimePoint startOfTracking = trackedRace.getStartOfTracking();
TimePoint endOfTracking = trackedRace.getEndOfTracking();
if (startOfTracking != null) {
GPSFix fixAfterStartOfTracking = markTrack
.getFirstFixAtOrAfter(startOfTracking);
if (fixAfterStartOfTracking == null
|| !trackedRace.isWithinStartAndEndOfTracking(
fixAfterStartOfTracking.getTimePoint())) {
// There is no fix in the tracking interval, so this fix could be "better"
// than ones already available in the track
// Better means closer before/after the beginning/end of the tracking
// interval
if (timePoint.before(startOfTracking)) {
// check if it is closer to the beginning of the tracking interval
GPSFix fixBeforeStartOfTracking = markTrack
.getLastFixAtOrBefore(startOfTracking);
forceFix = (fixBeforeStartOfTracking == null
|| fixBeforeStartOfTracking.getTimePoint().before(timePoint));
} else if (endOfTracking != null && timePoint.after(endOfTracking)) {
// check if it is closer to the end of the tracking interval
GPSFix fixAfterEndOfTracking = markTrack
.getFirstFixAtOrAfter(endOfTracking);
forceFix = (fixAfterEndOfTracking == null
|| fixAfterEndOfTracking.getTimePoint().after(timePoint));
} else {
forceFix = false;
}
} else {
// there is already a fix in the tracking interval
forceFix = false;
}
} else {
forceFix = false;
}
}
} finally {
markTrack.unlockAfterRead();
}
}
trackedRace.recordFix(mark, (GPSFix) fix, /* only when in tracking interval */ !forceFix);
if (returnLiveDelay) {
delayToLive.put(trackedRace.getRaceIdentifier(), new MillisecondsDurationImpl(trackedRace.getDelayToLiveInMillis()));
}
}
}
}
@Override
public void visit(RegattaLogDeviceCompetitorMappingEvent event) {
recordForCompetitor(event.getMappedTo());
}
@Override
public void visit(RegattaLogDeviceBoatMappingEvent event) {
final Boat boat = event.getMappedTo();
final Competitor comp = trackedRace.getCompetitorOfBoat(boat);
if (comp != null) {
recordForCompetitor(comp);
} else {
// this is not necessarily something to warn of; while a boat tracker may continuously track
logger.log(Level.FINE,
()->"Could not record fix for boat because no competitor could be determined. Boat: " + boat);
}
}
private void recordForCompetitor(Competitor comp) {
if (!preemptiveStopRequested.get()) {
if (fix instanceof GPSFixMoving) {
// try to record the fix, and only if it was really to the track,
// check for maneuvers; otherwise, the fix may not have been accepted
// by the race or the track, e.g., because the race's end-of-tracking
// comes before the fix's time point
if (trackedRace.recordFix(comp, (GPSFixMoving) fix)) { // TOOD bug6229: this checks the TrackedRace's tracking interval, but for an MDI we'd also want to intersect with Event/Regatta end date if set
if (returnManeuverChanges) {
RegattaAndRaceIdentifier maneuverChangedAnswer = detectIfManeuverChanged(comp);
if (maneuverChangedAnswer != null) {
maneuverChanged.add(maneuverChangedAnswer);
}
}
if (returnLiveDelay) {
delayToLive.put(trackedRace.getRaceIdentifier(), new MillisecondsDurationImpl(trackedRace.getDelayToLiveInMillis()));
}
}
} else {
logger.log(Level.WARNING,
String.format(
"Could not add fix for competitor (%s) in race (%s), as it"
+ " is no GPSFixMoving, meaning it is missing COG/SOG values",
comp, trackedRace.getRace().getName()));
}
}
}
@Override
public void visit(RegattaLogDeviceMarkMappingEvent event) {
if (!preemptiveStopRequested.get()) {
Mark mark = event.getMappedTo();
final DynamicGPSFixTrack<Mark, GPSFix> markTrack = trackedRace.getOrCreateTrack(mark);
final GPSFix firstFixAtOrAfter;
final boolean forceFix;
if (trackedRace.isWithinStartAndEndOfTracking(fix.getTimePoint())) {
forceFix = false;
} else {
markTrack.lockForRead();
try {
if (Util.isEmpty(markTrack.getRawFixes())
|| (firstFixAtOrAfter = markTrack.getFirstFixAtOrAfter(timePoint)) != null
&& firstFixAtOrAfter.getTimePoint().equals(timePoint)) {
// either the first fix or overwriting an existing one
forceFix = true;
} else {
// checking if the given fix is "better" than an existing one
TimePoint startOfTracking = trackedRace.getStartOfTracking();
TimePoint endOfTracking = trackedRace.getEndOfTracking();
if (startOfTracking != null) {
GPSFix fixAfterStartOfTracking = markTrack
.getFirstFixAtOrAfter(startOfTracking);
if (fixAfterStartOfTracking == null
|| !trackedRace.isWithinStartAndEndOfTracking(
fixAfterStartOfTracking.getTimePoint())) {
// There is no fix in the tracking interval, so this fix could be "better"
// than ones already available in the track
// Better means closer before/after the beginning/end of the tracking
// interval
if (timePoint.before(startOfTracking)) {
// check if it is closer to the beginning of the tracking interval
GPSFix fixBeforeStartOfTracking = markTrack
.getLastFixAtOrBefore(startOfTracking);
forceFix = (fixBeforeStartOfTracking == null
|| fixBeforeStartOfTracking.getTimePoint().before(timePoint));
} else if (endOfTracking != null && timePoint.after(endOfTracking)) {
// check if it is closer to the end of the tracking interval
GPSFix fixAfterEndOfTracking = markTrack
.getFirstFixAtOrAfter(endOfTracking);
forceFix = (fixAfterEndOfTracking == null
|| fixAfterEndOfTracking.getTimePoint().after(timePoint));
} else {
forceFix = false;
}
} else {
// there is already a fix in the tracking interval
forceFix = false;
}
} else {
forceFix = false;
}
}
} finally {
markTrack.unlockAfterRead();
}
}
trackedRace.recordFix(mark, (GPSFix) fix, /* only when in tracking interval */ !forceFix);
if (returnLiveDelay) {
delayToLive.put(trackedRace.getRaceIdentifier(), new MillisecondsDurationImpl(trackedRace.getDelayToLiveInMillis()));
}
}
}
});
}
});
});
}
});
}
}
return mergeManeuverChangedAndLiveDelayResult(maneuverChanged, delayToLive);
}
@@ -879,6 +880,11 @@ public class FixLoaderAndTracker implements TrackingDataLoader {
addLoadingJob(new LoadFixesForNewlyCoveredTimeRangesJob(item, newlyCoveredTimeRanges));
}
}
@Override
public String toString() {
return "FixLoaderDeviceMappings for race "+trackedRace.getRaceIdentifier();
}
}
/**
@@ -4,6 +4,7 @@ import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -69,6 +70,31 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
private final Map<ItemT, List<DeviceMappingWithRegattaLogEvent<ItemT>>> mappings = new HashMap<>();
private final Map<DeviceIdentifier, List<DeviceMappingWithRegattaLogEvent<ItemT>>> mappingsByDevice = new HashMap<>();
/**
* A cache that holds the device mappings as the {@link Pair#getB() second} component of the values in this map,
* such that exactly these device mappings apply for any time point {@link TimeRange#includes(TimePoint) included}
* by the {@link TimeRange} that is the {@link Pair#getA() first} component of a value in this map. This map's keys
* match with the {@link DeviceMapping#getDevice() device identifiers} of the {@link Pair#getB() second} components
* of their corresponding values.
* <p>
*
* This cache is designed to work well for cases where mappings change at a frequency orders of magnitude less than
* the frequency with which fixes arrive and are to be mapped to items. Furthermore, the cache hit rates benefit
* from mappings covering large time ranges.
* <p>
*
* Any change to the mappings for a device will remove the mapping for the device's {@link DeviceIdentifier
* identifier} from this map.
* <p>
*
* Access to this map has to undergo the same locking drill as any access to {@link #mappings}, using the
* {@link #mappingsLock}.
*/
private final Map<DeviceIdentifier, Pair<TimeRange, List<DeviceMappingWithRegattaLogEvent<ItemT>>>> cachedMappings = new HashMap<>();
private int cacheHits;
private int cacheMisses;
private final RegattaLogEventVisitor regattaLogEventVisitor = new BaseRegattaLogEventVisitor() {
@Override
public void visit(RegattaLogDeviceCompetitorSensorDataMappingEvent event) {
@@ -169,13 +195,38 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
public void forEachMappingOfDeviceIncludingTimePoint(DeviceIdentifier device, TimePoint timePoint,
Consumer<DeviceMappingWithRegattaLogEvent<ItemT>> callback) {
LockUtil.executeWithReadLock(mappingsLock, () -> {
List<DeviceMappingWithRegattaLogEvent<ItemT>> mappingsForDevice = mappingsByDevice.get(device);
if (mappingsForDevice != null) {
for (DeviceMappingWithRegattaLogEvent<ItemT> mapping : mappingsForDevice) {
if (mapping.getTimeRange().includes(timePoint)) {
callback.accept(mapping);
final Pair<TimeRange, List<DeviceMappingWithRegattaLogEvent<ItemT>>> cachedTimeRangeForDevice = cachedMappings.get(device);
if (cachedTimeRangeForDevice != null && cachedTimeRangeForDevice.getA().includes(timePoint)) {
cacheHits++;
logger.fine(() -> "Device mapping cache hit for mapper " + this + " for device " + device
+ " and time point " + timePoint + ", included in cached range "
+ cachedTimeRangeForDevice.getA() + "; " + cacheHits + " hits, " + cacheMisses + " misses");
cachedTimeRangeForDevice.getB().forEach(mapping->callback.accept(mapping));
} else {
final List<DeviceMappingWithRegattaLogEvent<ItemT>> mappingsForDevice = mappingsByDevice.get(device);
MultiTimeRange timeRangeForCache = MultiTimeRange.of();
final List<DeviceMappingWithRegattaLogEvent<ItemT>> deviceMappingsForCache = new LinkedList<>();
assert timeRangeForCache.isEmpty();
cacheMisses++;
if (mappingsForDevice != null) {
for (final DeviceMappingWithRegattaLogEvent<ItemT> mapping : mappingsForDevice) {
if (mapping.getTimeRange().includes(timePoint)) {
if (timeRangeForCache.isEmpty()) {
timeRangeForCache.union(mapping.getTimeRange());
} else {
timeRangeForCache = timeRangeForCache.intersection(mapping.getTimeRange());
}
deviceMappingsForCache.add(mapping);
callback.accept(mapping);
}
}
}
final MultiTimeRange finalTimeRangeForache = timeRangeForCache;
logger.fine(() -> "Device mapping cache miss for mapper " + this + " for device " + device
+ " and time point " + timePoint + ", determined cachable range "
+ finalTimeRangeForache + "; " + cacheHits + " hits, " + cacheMisses + " misses");
// TODO bug6236 Now issue the caching of the entry (device, (timeRangeForCache, deviceMappingsForCache)) under the mappingsLock's write lock!
}
});
}
@@ -245,6 +296,7 @@ public abstract class RegattaLogDeviceMappings<ItemT extends WithID> {
oldMappings.putAll(mappings);
oldDeviceIds.addAll(mappingsByDevice.keySet());
mappings.clear();
cachedMappings.clear();
mappings.putAll(newMappings);
mappingsByDevice.clear();
for (ItemT item : newMappings.keySet()) {