From f8b290844ba3b7696bf7ec4f8e6f75652ee6a09b Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 19 Mar 2021 10:16:33 -0400 Subject: [PATCH] GH-1046: Fix Out of Order Offset Commit with DLQ Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1046 When using ack mode `MANUAL_IMMEDIATE`, we must acknowledge the delivery on the container thread, otherwise the commit will be queued for later processing, possibly causing out of order commits, incorrectly reducing the committed offset. Wait for the send to complete before acknowledging; use a timeout slightly larger than the configured producer delivery timeout to avoid premature timeouts. Tested with user-provided test case. Fix missing timeout buffer. More context for this issue: https://gitter.im/spring-cloud/spring-cloud-stream?at=6050f43595e23446e43cacd1 --- .../kafka/KafkaMessageChannelBinder.java | 43 +++++++++++++++---- 1 file changed, 34 insertions(+), 9 deletions(-) diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 7882ab9c4..dc32c0a6f 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -30,6 +30,7 @@ import java.util.Map; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Predicate; import java.util.regex.Pattern; @@ -1141,8 +1142,21 @@ public class KafkaMessageChannelBinder extends final KafkaTemplate kafkaTemplate = new KafkaTemplate<>( producerFactory); + Object timeout = producerFactory.getConfigurationProperties().get(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG); + Long sendTimeout = null; + if (timeout instanceof Number) { + sendTimeout = ((Number) timeout).longValue() + 2000L; + } + else if (timeout instanceof String) { + sendTimeout = Long.parseLong((String) timeout) + 2000L; + } + if (timeout == null) { + sendTimeout = ((Integer) ProducerConfig.configDef() + .defaultValues() + .get(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG)).longValue() + 2000L; + } @SuppressWarnings("rawtypes") - DlqSender dlqSender = new DlqSender(kafkaTemplate); + DlqSender dlqSender = new DlqSender(kafkaTemplate, sendTimeout); return (message) -> { @@ -1574,8 +1588,11 @@ public class KafkaMessageChannelBinder extends private final KafkaTemplate kafkaTemplate; - DlqSender(KafkaTemplate kafkaTemplate) { + private final long sendTimeout; + + DlqSender(KafkaTemplate kafkaTemplate, long timeout) { this.kafkaTemplate = kafkaTemplate; + this.sendTimeout = timeout; } @SuppressWarnings("unchecked") @@ -1609,18 +1626,26 @@ public class KafkaMessageChannelBinder extends public void onSuccess(SendResult result) { if (KafkaMessageChannelBinder.this.logger.isDebugEnabled()) { KafkaMessageChannelBinder.this.logger - .debug("Sent to DLQ " + sb.toString()); - } - if (ackMode == ContainerProperties.AckMode.MANUAL || ackMode == ContainerProperties.AckMode.MANUAL_IMMEDIATE) { - messageHeaders.get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class).acknowledge(); + .debug("Sent to DLQ " + sb.toString() + ": " + result.getRecordMetadata()); } } }); + try { + sentDlq.get(this.sendTimeout, TimeUnit.MILLISECONDS); + } + catch (InterruptedException ex) { + Thread.currentThread().interrupt(); + throw ex; + } } catch (Exception ex) { - if (sentDlq == null) { - KafkaMessageChannelBinder.this.logger - .error("Error sending to DLQ " + sb.toString(), ex); + KafkaMessageChannelBinder.this.logger + .error("Error sending to DLQ " + sb.toString(), ex); + } + finally { + if (ackMode == ContainerProperties.AckMode.MANUAL + || ackMode == ContainerProperties.AckMode.MANUAL_IMMEDIATE) { + messageHeaders.get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class).acknowledge(); } }