From 8d5865b64cd44b6c1475b8240e8989c33e00d349 Mon Sep 17 00:00:00 2001 From: Alexander Rovner Date: Tue, 3 Mar 2026 12:45:41 +0100 Subject: [PATCH 1/2] test whether it is possible to consume before the producer gets an ack spoiler: yes it's possible --- src/main/java/io/spoud/MetricService.java | 19 +++++++++++++++---- .../java/io/spoud/kafka/MessageConsumer.java | 2 +- 2 files changed, 16 insertions(+), 5 deletions(-) diff --git a/src/main/java/io/spoud/MetricService.java b/src/main/java/io/spoud/MetricService.java index fd3e407..6b38e66 100644 --- a/src/main/java/io/spoud/MetricService.java +++ b/src/main/java/io/spoud/MetricService.java @@ -4,6 +4,7 @@ import io.quarkus.logging.Log; import io.spoud.config.SynthClientConfig; import io.spoud.kafka.PartitionRebalancer; +import io.vertx.core.impl.ConcurrentHashSet; import jakarta.enterprise.context.ApplicationScoped; import org.eclipse.microprofile.config.inject.ConfigProperty; @@ -89,7 +90,12 @@ public void recordProducedFailure() { recordsFailedCounter.increment(); } - synchronized public void recordLatency(String topic, int partition, long latencyMs, String fromRack) { + synchronized public void recordLatency(String topic, int partition, long latencyMs, String fromRack, long offset) { + if (fromRack.equals(config.rack()) && !ackedOffsets.contains(new AckedMessage(offset, partition))) { + Log.warnf("Consumed message %d:%d before receiving an ACK! This means that Ack latency could be higher than E2E latency!", partition, offset); + } else { + ackedOffsets.remove(new AckedMessage(offset, partition)); + } Log.debugv("Latency for partition {0}: {1}ms", partition, latencyMs); if (partitionRebalancer.isInitialRefreshPending()) { Log.info("Ignoring latencies as the initial partition assignment is not done yet"); @@ -132,7 +138,12 @@ public Collection getE2ELatencies() { return e2eLatencies.values(); } - public void recordAckLatency(String topic, int partition, Duration between) { + public record AckedMessage(long offset, int partition) { + } + private final Set ackedOffsets = new ConcurrentHashSet<>(); + + public void recordAckLatency(String topic, int partition, Duration between, long offset) { + ackedOffsets.add(new AckedMessage(offset, partition)); Log.debugv("Ack latency for partition {0}: {1}ms", partition, between.toMillis()); if (partitionRebalancer.isInitialRefreshPending()) { Log.info("Ignoring ack latency as the initial partition assignment is not done yet"); @@ -192,8 +203,8 @@ private WrappedDistributionSummary genE2eSummary(String topic, int partition, St .tag(TAG_FROM_RACK, fromRack) .tag(TAG_BROKER_RACK, brokerRack) .description("End-to-end latency of the synthetic client") - .minimumExpectedValue(1.0) - .maximumExpectedValue(10_000.0) + .minimumExpectedValue(config.expectedMinLatency()) + .maximumExpectedValue(config.expectedMaxLatency()) .publishPercentiles(0.5, 0.8, 0.9, 0.95, 0.99) .publishPercentileHistogram(config.publishHistogramBuckets()) .distributionStatisticExpiry(config.samplingTimeWindow()) diff --git a/src/main/java/io/spoud/kafka/MessageConsumer.java b/src/main/java/io/spoud/kafka/MessageConsumer.java index 53da40b..ad7b141 100644 --- a/src/main/java/io/spoud/kafka/MessageConsumer.java +++ b/src/main/java/io/spoud/kafka/MessageConsumer.java @@ -108,7 +108,7 @@ public void onPartitionsAssigned(Collection collection) { .map(String::new) .orElse(null); metricService.recordConsumptionTime(); - metricService.recordLatency(message.topic(), message.partition(), consumeTime - produceTime, fromRack); + metricService.recordLatency(message.topic(), message.partition(), consumeTime - produceTime, fromRack, message.offset()); advertisedListenerRepository.mapRackToUrl(fromRack, advertisedListener); lastReport.updateAndGet(last -> { if (Duration.between(last, Instant.now()).getSeconds() > 10) { From bc27353bff02e6fdf7066efe06e39862699ac72d Mon Sep 17 00:00:00 2001 From: Alexander Rovner Date: Tue, 3 Mar 2026 13:02:41 +0100 Subject: [PATCH 2/2] more consistent timing --- src/main/java/io/spoud/kafka/MessageProducer.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/src/main/java/io/spoud/kafka/MessageProducer.java b/src/main/java/io/spoud/kafka/MessageProducer.java index d5e63be..a51f9ac 100644 --- a/src/main/java/io/spoud/kafka/MessageProducer.java +++ b/src/main/java/io/spoud/kafka/MessageProducer.java @@ -78,7 +78,6 @@ public void recreateProducer() { } public void send(Long key, String value) { - Instant send = Instant.now(); var record = new ProducerRecord<>(config.topic(), null, timeService.currentTimeMillis(), key, value); record.headers().add(HEADER_RACK, config.rack().getBytes()); record.headers().add(HEADER_ADVERTISED_LISTENER, config.advertisedListener().orElse("").getBytes()); @@ -87,9 +86,10 @@ var record = new ProducerRecord<>(config.topic(), null, timeService.currentTimeM Log.error("Failed to send message", exception); metricService.recordProducedFailure(); } else { - Instant ack = Instant.now(); - lastMessage.set(ack); - metricService.recordAckLatency(metadata.topic(), metadata.partition(), Duration.between(send, ack)); + Instant ackTime = Instant.ofEpochMilli(timeService.currentTimeMillis()); + Instant sendTime = Instant.ofEpochMilli(metadata.timestamp()); + lastMessage.set(ackTime); + metricService.recordAckLatency(metadata.topic(), metadata.partition(), Duration.between(sendTime, ackTime), metadata.offset()); metricService.recordProducedSuccess(); } });