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
This commit is contained in:
Gary Russell
2017-12-28 15:59:00 -05:00
committed by Artem Bilan
parent fb9dd91644
commit 29973e525d
2 changed files with 28 additions and 3 deletions

View File

@@ -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<K, V> 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<K, V> 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 <T, E extends Throwable> boolean open(RetryContext context, RetryCallback<T, E> callback) {
if (KafkaMessageDrivenChannelAdapter.this.recoveryCallback != null) {

View File

@@ -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();
}