From bdd236052b61e91581170ac978fd2ae4dee9db28 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 17 Mar 2020 12:09:21 -0400 Subject: [PATCH] Upgrade to SK 2.5.0 - add deliveryAttempt header when retry is not configured - temporary work around for gradle bug, jar with test classifier missing from CP --- .../integration/kafka/dsl/Kafka.java | 48 ------------------- .../kafka/inbound/KafkaInboundGateway.java | 18 ++++++- .../KafkaMessageDrivenChannelAdapter.java | 18 ++++++- .../inbound/MessageDrivenAdapterTests.java | 6 ++- 4 files changed, 39 insertions(+), 51 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java index b03aefa79b..b35592c4e2 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/dsl/Kafka.java @@ -16,7 +16,6 @@ package org.springframework.integration.kafka.dsl; -import java.util.Arrays; import java.util.regex.Pattern; import org.apache.kafka.common.TopicPartition; @@ -286,26 +285,6 @@ public final class Kafka { containerProperties), listenerMode); } - /** - * Create an initial - * {@link KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec}. - * @param consumerFactory the {@link ConsumerFactory}. - * @param topicPartitions the {@link TopicPartition} vararg. - * @param the Kafka message key type. - * @param the Kafka message value type. - * @return the KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec. - * @deprecated in favor of {@link #messageDrivenChannelAdapter(ConsumerFactory, TopicPartitionOffset...)}. - */ - @Deprecated - public static - KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec messageDrivenChannelAdapter( - ConsumerFactory consumerFactory, - org.springframework.kafka.support.TopicPartitionInitialOffset... topicPartitions) { - - return messageDrivenChannelAdapter(consumerFactory, KafkaMessageDrivenChannelAdapter.ListenerMode.record, - topicPartitions); - } - /** * Create an initial * {@link KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec}. @@ -324,33 +303,6 @@ public final class Kafka { topicPartitions); } - /** - * Create an initial - * {@link KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec}. - * @param consumerFactory the {@link ConsumerFactory}. - * @param listenerMode the {@link KafkaMessageDrivenChannelAdapter.ListenerMode}. - * @param topicPartitions the {@link TopicPartition} vararg. - * @param the Kafka message key type. - * @param the Kafka message value type. - * @return the - * KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec. - * @deprecated in favor of - * {@link #messageDrivenChannelAdapter(ConsumerFactory, org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter.ListenerMode, TopicPartitionOffset...)} - */ - @Deprecated - public static - KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec messageDrivenChannelAdapter( - ConsumerFactory consumerFactory, - KafkaMessageDrivenChannelAdapter.ListenerMode listenerMode, - org.springframework.kafka.support.TopicPartitionInitialOffset... topicPartitions) { - - return messageDrivenChannelAdapter( - new KafkaMessageListenerContainerSpec<>(consumerFactory, - Arrays.stream(topicPartitions) - .map(org.springframework.kafka.support.TopicPartitionInitialOffset::toTPO) - .toArray(TopicPartitionOffset[]::new)), listenerMode); - } - /** * Create an initial * {@link KafkaMessageDrivenChannelAdapterSpec.KafkaMessageDrivenChannelAdapterListenerContainerSpec}. diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java index 448259ed9d..9d51c1ef1d 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java @@ -1,5 +1,5 @@ /* - * Copyright 2018-2019 the original author or authors. + * Copyright 2018-2020 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. @@ -16,6 +16,7 @@ package org.springframework.integration.kafka.inbound; +import java.nio.ByteBuffer; import java.util.Map; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BiConsumer; @@ -23,6 +24,7 @@ import java.util.function.BiConsumer; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.header.Header; import org.springframework.core.AttributeAccessor; import org.springframework.integration.IntegrationMessageHeaderAccessor; @@ -86,6 +88,8 @@ public class KafkaInboundGateway extends MessagingGatewaySupport implem private boolean bindSourceRecord; + private boolean containerDeliveryAttemptPresent; + /** * Construct an instance with the provided container. * @param messageListenerContainer the container. @@ -178,6 +182,8 @@ public class KafkaInboundGateway extends MessagingGatewaySupport implem this.retryTemplate.registerListener(this.listener); } this.messageListenerContainer.getContainerProperties().setMessageListener(kafkaListener); + this.containerDeliveryAttemptPresent = this.messageListenerContainer.getContainerProperties() + .isDeliveryAttemptHeader(); } @Override @@ -295,6 +301,11 @@ public class KafkaInboundGateway extends MessagingGatewaySupport implem new AtomicInteger(((RetryContext) attributesHolder.get()).getRetryCount() + 1); rawHeaders.put(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, deliveryAttempt); } + else if (KafkaInboundGateway.this.containerDeliveryAttemptPresent) { + Header header = record.headers().lastHeader(KafkaHeaders.DELIVERY_ATTEMPT); + rawHeaders.put(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, + new AtomicInteger(ByteBuffer.wrap(header.value()).getInt())); + } if (KafkaInboundGateway.this.bindSourceRecord) { rawHeaders.put(IntegrationMessageHeaderAccessor.SOURCE_DATA, record); } @@ -306,6 +317,11 @@ public class KafkaInboundGateway extends MessagingGatewaySupport implem new AtomicInteger(((RetryContext) attributesHolder.get()).getRetryCount() + 1); builder.setHeader(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, deliveryAttempt); } + else if (KafkaInboundGateway.this.containerDeliveryAttemptPresent) { + Header header = record.headers().lastHeader(KafkaHeaders.DELIVERY_ATTEMPT); + builder.setHeader(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, + new AtomicInteger(ByteBuffer.wrap(header.value()).getInt())); + } if (KafkaInboundGateway.this.bindSourceRecord) { builder.setHeader(IntegrationMessageHeaderAccessor.SOURCE_DATA, record); } 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 d35da81c7d..7d2255dfa4 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-2019 the original author or authors. + * Copyright 2015-2020 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. @@ -16,6 +16,7 @@ package org.springframework.integration.kafka.inbound; +import java.nio.ByteBuffer; import java.util.List; import java.util.Map; import java.util.concurrent.atomic.AtomicInteger; @@ -24,6 +25,7 @@ import java.util.function.BiConsumer; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.header.Header; import org.springframework.core.AttributeAccessor; import org.springframework.integration.IntegrationMessageHeaderAccessor; @@ -99,6 +101,8 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo private boolean bindSourceRecord; + private boolean containerDeliveryAttemptPresent; + /** * Construct an instance with mode {@link ListenerMode#record}. * @param messageListenerContainer the container. @@ -310,6 +314,8 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo } this.messageListenerContainer.getContainerProperties().setMessageListener(listener); } + this.containerDeliveryAttemptPresent = this.messageListenerContainer.getContainerProperties() + .isDeliveryAttemptHeader(); } @Override @@ -452,6 +458,11 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo new AtomicInteger(((RetryContext) attributesHolder.get()).getRetryCount() + 1); rawHeaders.put(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, deliveryAttempt); } + else if (KafkaMessageDrivenChannelAdapter.this.containerDeliveryAttemptPresent) { + Header header = record.headers().lastHeader(KafkaHeaders.DELIVERY_ATTEMPT); + rawHeaders.put(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, + new AtomicInteger(ByteBuffer.wrap(header.value()).getInt())); + } if (KafkaMessageDrivenChannelAdapter.this.bindSourceRecord) { rawHeaders.put(IntegrationMessageHeaderAccessor.SOURCE_DATA, record); } @@ -463,6 +474,11 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo new AtomicInteger(((RetryContext) attributesHolder.get()).getRetryCount() + 1); builder.setHeader(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, deliveryAttempt); } + else if (KafkaMessageDrivenChannelAdapter.this.containerDeliveryAttemptPresent) { + Header header = record.headers().lastHeader(KafkaHeaders.DELIVERY_ATTEMPT); + builder.setHeader(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, + new AtomicInteger(ByteBuffer.wrap(header.value()).getInt())); + } if (KafkaMessageDrivenChannelAdapter.this.bindSourceRecord) { builder.setHeader(IntegrationMessageHeaderAccessor.SOURCE_DATA, record); } 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 1fc76585d4..208122f0af 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-2019 the original author or authors. + * Copyright 2016-2020 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. @@ -66,6 +66,7 @@ import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.listener.KafkaMessageListenerContainer; +import org.springframework.kafka.listener.SeekToCurrentErrorHandler; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.DefaultKafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaders; @@ -329,8 +330,10 @@ class MessageDrivenAdapterTests { props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); DefaultKafkaConsumerFactory cf = new DefaultKafkaConsumerFactory<>(props); ContainerProperties containerProps = new ContainerProperties(topic5); + containerProps.setDeliveryAttemptHeader(true); KafkaMessageListenerContainer container = new KafkaMessageListenerContainer<>(cf, containerProps); + container.setErrorHandler(new SeekToCurrentErrorHandler()); KafkaMessageDrivenChannelAdapter adapter = new KafkaMessageDrivenChannelAdapter<>(container); MessageChannel out = new DirectChannel() { @@ -371,6 +374,7 @@ class MessageDrivenAdapterTests { assertThat(headers.get(KafkaHeaders.RECEIVED_TOPIC)).isEqualTo(topic5); assertThat(headers.get(KafkaHeaders.RECEIVED_PARTITION_ID)).isEqualTo(0); assertThat(headers.get(KafkaHeaders.OFFSET)).isEqualTo(0L); + assertThat(StaticMessageHeaderAccessor.getDeliveryAttempt(originalMessage).get()).isEqualTo(1); adapter.stop(); }