From 29973e525de5ea2686cf08e1f487e4e1bdda7adf Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 28 Dec 2017 15:59:00 -0500 Subject: [PATCH] GH-185: Add deliveryAttempts header with retry Resolves: https://github.com/spring-projects/spring-integration-kafka/issues/185 When a `RetryTemplate` is wired into the message-driven adapter, add and increment an `AtomicInteger`-valued header. See https://jira.spring.io/browse/INT-4369 * Actualize Copyright in the affected classes --- .../KafkaMessageDrivenChannelAdapter.java | 25 ++++++++++++++++++- .../inbound/MessageDrivenAdapterTests.java | 6 +++-- 2 files changed, 28 insertions(+), 3 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java index 948aa13531..6f13b7945c 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2015-2017 the original author or authors. + * Copyright 2015-2018 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -17,16 +17,19 @@ package org.springframework.integration.kafka.inbound; import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.core.AttributeAccessor; +import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.context.OrderlyShutdownCapable; import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.integration.kafka.support.RawRecordHeaderErrorMessageStrategy; import org.springframework.integration.support.ErrorMessageStrategy; import org.springframework.integration.support.ErrorMessageUtils; +import org.springframework.integration.support.MessageBuilder; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.BatchMessageListener; import org.springframework.kafka.listener.MessageListener; @@ -40,6 +43,7 @@ import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.converter.BatchMessageConverter; import org.springframework.kafka.support.converter.ConversionException; +import org.springframework.kafka.support.converter.KafkaMessageHeaders; import org.springframework.kafka.support.converter.MessageConverter; import org.springframework.kafka.support.converter.RecordMessageConverter; import org.springframework.messaging.Message; @@ -358,6 +362,9 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo Message message = null; try { message = toMessagingMessage(record, acknowledgment, consumer); + if (KafkaMessageDrivenChannelAdapter.this.retryTemplate != null) { + message = addDeliveryAttemptHeader(message); + } setAttributesIfNecessary(record, message); } catch (RuntimeException e) { @@ -380,6 +387,22 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo } } + private Message addDeliveryAttemptHeader(Message message) { + Message messageToReturn = message; + AtomicInteger deliveryAttempt = + new AtomicInteger(((RetryContext) attributesHolder.get()).getRetryCount() + 1); + if (message.getHeaders() instanceof KafkaMessageHeaders) { + ((KafkaMessageHeaders) message.getHeaders()).getRawHeaders() + .put(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, deliveryAttempt); + } + else { + messageToReturn = MessageBuilder.fromMessage(message) + .setHeader(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, deliveryAttempt) + .build(); + } + return messageToReturn; + } + @Override public boolean open(RetryContext context, RetryCallback callback) { if (KafkaMessageDrivenChannelAdapter.this.recoveryCallback != null) { diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java index 88cb6bfd4a..6435c149e6 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/inbound/MessageDrivenAdapterTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2017 the original author or authors. + * Copyright 2016-2018 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -39,6 +39,7 @@ import org.springframework.integration.handler.advice.ErrorMessageSendingRecover import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.ListenerMode; import org.springframework.integration.kafka.support.RawRecordHeaderErrorMessageStrategy; import org.springframework.integration.support.MessageBuilder; +import org.springframework.integration.support.StaticMessageHeaderAccessor; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.core.KafkaTemplate; @@ -195,7 +196,7 @@ public class MessageDrivenAdapterTests { adapter.setOutputChannel(out); RetryTemplate retryTemplate = new RetryTemplate(); SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); - retryPolicy.setMaxAttempts(1); + retryPolicy.setMaxAttempts(2); retryTemplate.setRetryPolicy(retryPolicy); QueueChannel errorChannel = new QueueChannel(); adapter.setRecoveryCallback( @@ -222,6 +223,7 @@ public class MessageDrivenAdapterTests { assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic4); assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0); assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L); + assertThat(StaticMessageHeaderAccessor.getDeliveryAttempt(received).get()).isEqualTo(2); adapter.stop(); }