Skip to content

Commit c7f4ba6

Browse files
Chance Newkirkclaude
authored andcommitted
Fix out-of-order sample loss by making remote write synchronous
The store() method previously fired HTTP writes asynchronously via executeAsync() and returned immediately. When the RingBuffer's multiple worker threads dispatched consecutive batches containing samples for the same series, the async HTTP requests could arrive at the remote write endpoint out of timestamp order, causing the backend to reject the stale samples as out-of-order. This change makes store() block until the HTTP write completes, ensuring the ring buffer worker thread does not process the next batch until the current write has landed. This preserves per-series timestamp ordering across consecutive WriteRequests as required by the Prometheus Remote Write spec. Additionally fixes a bug where samplesLost incorrectly counted unfiltered samples (including NaN) instead of the actual samples that were attempted. Validated via A/B E2E testing against Thanos Receive: - Baseline (async): 14 out-of-order / 5,262 appended (0.27% loss) - Fix (sync): 0 out-of-order / 5,264 appended (0.00% loss) - Throughput: identical (~5,260 samples over equal soak periods) - All 45 smoke tests passing on both runs Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent fe9bda4 commit c7f4ba6

1 file changed

Lines changed: 13 additions & 9 deletions

File tree

‎plugin/src/main/java/org/opennms/timeseries/cortex/CortexTSS.java‎

Lines changed: 13 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -233,15 +233,19 @@ public void store(final List<Sample> samples, String clientID) throws StorageExc
233233
final Request request = builder.build();
234234

235235
LOG.trace("Writing: {}", writeRequest);
236-
asyncHttpCallsBulkhead.executeCompletionStage(() -> executeAsync(request)).whenComplete((r, ex) -> {
237-
if (ex == null) {
238-
samplesWritten.mark(samplesSorted.size());
239-
} else {
240-
// FIXME: Data loss
241-
samplesLost.mark(samples.size());
242-
LOG.error("Error occurred while storing samples, sample will be lost.", ex);
243-
}
244-
});
236+
// Block until the HTTP write completes. This ensures the ring buffer worker
237+
// thread does not process the next batch until this write has landed, preserving
238+
// per-series timestamp ordering across consecutive WriteRequests as required by
239+
// the Prometheus Remote Write spec.
240+
try {
241+
asyncHttpCallsBulkhead.executeCompletionStage(() -> executeAsync(request)).toCompletableFuture().get(
242+
config.getWriteTimeoutInMs(), TimeUnit.MILLISECONDS);
243+
samplesWritten.mark(samplesSorted.size());
244+
} catch (Exception ex) {
245+
samplesLost.mark(samplesSorted.size());
246+
Throwable cause = ex.getCause() != null ? ex.getCause() : ex;
247+
throw new StorageException("Failed to write samples to Prometheus: " + cause.getMessage(), cause);
248+
}
245249
}
246250

247251
private void persistExternalTags(final Sample s) {

0 commit comments

Comments
 (0)