GH-164: Spring Kafka 2.0.0 Compatibility
Resolves: https://github.com/spring-projects/spring-integration-kafka/issues/164 Also gradle 3.5. Requires https://github.com/spring-projects/spring-kafka/pull/296 Fix javadocs More javadoc polishing Updates for new Consumer header * Simple polishing
This commit is contained in:
committed by
Artem Bilan
parent
69e533f338
commit
ee6b3f87b3
@@ -33,9 +33,7 @@ import org.springframework.kafka.listener.AbstractMessageListenerContainer;
|
||||
import org.springframework.kafka.listener.AcknowledgingMessageListener;
|
||||
import org.springframework.kafka.listener.ConcurrentMessageListenerContainer;
|
||||
import org.springframework.kafka.listener.ErrorHandler;
|
||||
import org.springframework.kafka.listener.adapter.FilteringAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.RecordFilterStrategy;
|
||||
import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.config.ContainerProperties;
|
||||
import org.springframework.kafka.support.TopicPartitionInitialOffset;
|
||||
import org.springframework.kafka.support.converter.BatchMessageConverter;
|
||||
@@ -99,7 +97,7 @@ public class KafkaMessageDrivenChannelAdapterSpec<K, V, S extends KafkaMessageDr
|
||||
/**
|
||||
* Specify a {@link RecordFilterStrategy} to wrap
|
||||
* {@code KafkaMessageDrivenChannelAdapter.IntegrationRecordMessageListener} into
|
||||
* {@link FilteringAcknowledgingMessageListenerAdapter}.
|
||||
* {@link org.springframework.kafka.listener.adapter.FilteringMessageListenerAdapter}.
|
||||
* @param recordFilterStrategy the {@link RecordFilterStrategy} to use.
|
||||
* @return the spec
|
||||
*/
|
||||
@@ -109,9 +107,10 @@ public class KafkaMessageDrivenChannelAdapterSpec<K, V, S extends KafkaMessageDr
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@code boolean} flag to indicate if {@link FilteringAcknowledgingMessageListenerAdapter}
|
||||
* should acknowledge discarded records or not.
|
||||
* Does not make sense if {@link #recordFilterStrategy(RecordFilterStrategy)} isn't specified.
|
||||
* A {@code boolean} flag to indicate if
|
||||
* {@link org.springframework.kafka.listener.adapter.FilteringMessageListenerAdapter}
|
||||
* should acknowledge discarded records or not. Does not make sense if
|
||||
* {@link #recordFilterStrategy(RecordFilterStrategy)} isn't specified.
|
||||
* @param ackDiscarded true to ack (commit offset for) discarded messages.
|
||||
* @return the spec
|
||||
*/
|
||||
@@ -123,7 +122,7 @@ public class KafkaMessageDrivenChannelAdapterSpec<K, V, S extends KafkaMessageDr
|
||||
/**
|
||||
* Specify a {@link RetryTemplate} instance to wrap
|
||||
* {@code KafkaMessageDrivenChannelAdapter.IntegrationRecordMessageListener} into
|
||||
* {@link RetryingAcknowledgingMessageListenerAdapter}.
|
||||
* {@link org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter}.
|
||||
* @param retryTemplate the {@link RetryTemplate} to use.
|
||||
* @return the spec
|
||||
*/
|
||||
@@ -145,15 +144,17 @@ public class KafkaMessageDrivenChannelAdapterSpec<K, V, S extends KafkaMessageDr
|
||||
}
|
||||
|
||||
/**
|
||||
/**
|
||||
* The {@code boolean} flag to specify the order how
|
||||
* {@link RetryingAcknowledgingMessageListenerAdapter} and
|
||||
* {@link FilteringAcknowledgingMessageListenerAdapter} are wrapped to each other,
|
||||
* if both of them are present.
|
||||
* Does not make sense if only one of {@link RetryTemplate} or
|
||||
* {@link RecordFilterStrategy} is present, or any.
|
||||
* @param filterInRetry the order for {@link RetryingAcknowledgingMessageListenerAdapter} and
|
||||
* {@link FilteringAcknowledgingMessageListenerAdapter} wrapping. Defaults to {@code false}.
|
||||
* {@link org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter}
|
||||
* and
|
||||
* {@link org.springframework.kafka.listener.adapter.FilteringMessageListenerAdapter}
|
||||
* are wrapped to each other, if both of them are present. Does not make sense if only
|
||||
* one of {@link RetryTemplate} or {@link RecordFilterStrategy} is present, or any.
|
||||
* @param filterInRetry the order for
|
||||
* {@link org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter}
|
||||
* and
|
||||
* {@link org.springframework.kafka.listener.adapter.FilteringMessageListenerAdapter}
|
||||
* wrapping. Defaults to {@code false}.
|
||||
* @return the spec
|
||||
*/
|
||||
public S filterInRetry(boolean filterInRetry) {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2015-2016 the original author or authors.
|
||||
* Copyright 2015-2017 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.
|
||||
@@ -18,19 +18,20 @@ package org.springframework.integration.kafka.inbound;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
|
||||
import org.springframework.integration.context.OrderlyShutdownCapable;
|
||||
import org.springframework.integration.endpoint.MessageProducerSupport;
|
||||
import org.springframework.kafka.listener.AbstractMessageListenerContainer;
|
||||
import org.springframework.kafka.listener.AcknowledgingMessageListener;
|
||||
import org.springframework.kafka.listener.BatchAcknowledgingMessageListener;
|
||||
import org.springframework.kafka.listener.BatchMessageListener;
|
||||
import org.springframework.kafka.listener.MessageListener;
|
||||
import org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.FilteringAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.FilteringBatchAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.FilteringBatchMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.FilteringMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.RecordFilterStrategy;
|
||||
import org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter;
|
||||
import org.springframework.kafka.support.Acknowledgment;
|
||||
import org.springframework.kafka.support.converter.BatchMessageConverter;
|
||||
import org.springframework.kafka.support.converter.ConversionException;
|
||||
@@ -69,7 +70,7 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
|
||||
|
||||
private RetryTemplate retryTemplate;
|
||||
|
||||
private RecoveryCallback<Void> recoveryCallback;
|
||||
private RecoveryCallback<? extends Object> recoveryCallback;
|
||||
|
||||
private boolean filterInRetry;
|
||||
|
||||
@@ -137,7 +138,7 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
|
||||
/**
|
||||
* Specify a {@link RecordFilterStrategy} to wrap
|
||||
* {@link KafkaMessageDrivenChannelAdapter.IntegrationRecordMessageListener} into
|
||||
* {@link FilteringAcknowledgingMessageListenerAdapter}.
|
||||
* {@link FilteringMessageListenerAdapter}.
|
||||
* @param recordFilterStrategy the {@link RecordFilterStrategy} to use.
|
||||
* @since 2.0.1
|
||||
*/
|
||||
@@ -146,7 +147,7 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@code boolean} flag to indicate if {@link FilteringAcknowledgingMessageListenerAdapter}
|
||||
* A {@code boolean} flag to indicate if {@link FilteringMessageListenerAdapter}
|
||||
* should acknowledge discarded records or not.
|
||||
* Does not make sense if {@link #setRecordFilterStrategy(RecordFilterStrategy)} isn't specified.
|
||||
* @param ackDiscarded true to ack (commit offset for) discarded messages.
|
||||
@@ -159,7 +160,7 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
|
||||
/**
|
||||
* Specify a {@link RetryTemplate} instance to wrap
|
||||
* {@link KafkaMessageDrivenChannelAdapter.IntegrationRecordMessageListener} into
|
||||
* {@link RetryingAcknowledgingMessageListenerAdapter}.
|
||||
* {@link RetryingMessageListenerAdapter}.
|
||||
* @param retryTemplate the {@link RetryTemplate} to use.
|
||||
* @since 2.0.1
|
||||
*/
|
||||
@@ -176,19 +177,19 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
|
||||
* @param recoveryCallback the recovery callback.
|
||||
* @since 2.0.1
|
||||
*/
|
||||
public void setRecoveryCallback(RecoveryCallback<Void> recoveryCallback) {
|
||||
public void setRecoveryCallback(RecoveryCallback<? extends Object> recoveryCallback) {
|
||||
this.recoveryCallback = recoveryCallback;
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@code boolean} flag to specify the order how
|
||||
* {@link RetryingAcknowledgingMessageListenerAdapter} and
|
||||
* {@link FilteringAcknowledgingMessageListenerAdapter} are wrapped to each other,
|
||||
* {@link RetryingMessageListenerAdapter} and
|
||||
* {@link FilteringMessageListenerAdapter} are wrapped to each other,
|
||||
* if both of them are present.
|
||||
* Does not make sense if only one of {@link RetryTemplate} or
|
||||
* {@link RecordFilterStrategy} is present, or any.
|
||||
* @param filterInRetry the order for {@link RetryingAcknowledgingMessageListenerAdapter} and
|
||||
* {@link FilteringAcknowledgingMessageListenerAdapter} wrapping. Defaults to {@code false}.
|
||||
* @param filterInRetry the order for {@link RetryingMessageListenerAdapter} and
|
||||
* {@link FilteringMessageListenerAdapter} wrapping. Defaults to {@code false}.
|
||||
* @since 2.0.1
|
||||
*/
|
||||
public void setFilterInRetry(boolean filterInRetry) {
|
||||
@@ -211,34 +212,34 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
|
||||
super.onInit();
|
||||
|
||||
if (this.mode.equals(ListenerMode.record)) {
|
||||
AcknowledgingMessageListener<K, V> listener = this.recordListener;
|
||||
MessageListener<K, V> listener = this.recordListener;
|
||||
|
||||
boolean filterInRetry = this.filterInRetry && this.retryTemplate != null
|
||||
&& this.recordFilterStrategy != null;
|
||||
|
||||
if (filterInRetry) {
|
||||
listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy,
|
||||
listener = new FilteringMessageListenerAdapter<>(listener, this.recordFilterStrategy,
|
||||
this.ackDiscarded);
|
||||
listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate,
|
||||
this.recoveryCallback);
|
||||
listener = new RetryingMessageListenerAdapter<>(listener, this.retryTemplate,
|
||||
this.recoveryCallback);
|
||||
}
|
||||
else {
|
||||
if (this.retryTemplate != null) {
|
||||
listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate,
|
||||
listener = new RetryingMessageListenerAdapter<>(listener, this.retryTemplate,
|
||||
this.recoveryCallback);
|
||||
}
|
||||
if (this.recordFilterStrategy != null) {
|
||||
listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy,
|
||||
listener = new FilteringMessageListenerAdapter<>(listener, this.recordFilterStrategy,
|
||||
this.ackDiscarded);
|
||||
}
|
||||
}
|
||||
this.messageListenerContainer.getContainerProperties().setMessageListener(listener);
|
||||
}
|
||||
else {
|
||||
BatchAcknowledgingMessageListener<K, V> listener = this.batchListener;
|
||||
BatchMessageListener<K, V> listener = this.batchListener;
|
||||
|
||||
if (this.recordFilterStrategy != null) {
|
||||
listener = new FilteringBatchAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy,
|
||||
listener = new FilteringBatchMessageListenerAdapter<>(listener, this.recordFilterStrategy,
|
||||
this.ackDiscarded);
|
||||
}
|
||||
this.messageListenerContainer.getContainerProperties().setMessageListener(listener);
|
||||
@@ -297,10 +298,10 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onMessage(ConsumerRecord<K, V> record, Acknowledgment acknowledgment) {
|
||||
public void onMessage(ConsumerRecord<K, V> record, Acknowledgment acknowledgment, Consumer<?, ?> consumer) {
|
||||
Message<?> message = null;
|
||||
try {
|
||||
message = toMessagingMessage(record, acknowledgment);
|
||||
message = toMessagingMessage(record, acknowledgment, consumer);
|
||||
}
|
||||
catch (RuntimeException e) {
|
||||
Exception exception = new ConversionException("Failed to convert to message for: " + record, e);
|
||||
@@ -326,10 +327,11 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onMessage(List<ConsumerRecord<K, V>> records, Acknowledgment acknowledgment) {
|
||||
Message<?> message = null;
|
||||
public void onMessage(List<ConsumerRecord<K, V>> records, Acknowledgment acknowledgment,
|
||||
Consumer<?, ?> consumer) {
|
||||
Message<?> message = null;
|
||||
try {
|
||||
message = toMessagingMessage(records, acknowledgment);
|
||||
message = toMessagingMessage(records, acknowledgment, consumer);
|
||||
}
|
||||
catch (RuntimeException e) {
|
||||
Exception exception = new ConversionException("Failed to convert to message for: " + records, e);
|
||||
|
||||
@@ -33,9 +33,9 @@ import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAd
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
|
||||
import org.springframework.kafka.listener.KafkaMessageListenerContainer;
|
||||
import org.springframework.kafka.listener.adapter.FilteringAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.FilteringMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.RecordFilterStrategy;
|
||||
import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.config.ContainerProperties;
|
||||
import org.springframework.retry.support.RetryTemplate;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
@@ -112,7 +112,7 @@ public class KafkaMessageDrivenChannelAdapterParserTests {
|
||||
containerProps = TestUtils.getPropertyValue(container, "containerProperties", ContainerProperties.class);
|
||||
|
||||
Object messageListener = containerProps.getMessageListener();
|
||||
assertThat(messageListener).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class);
|
||||
assertThat(messageListener).isInstanceOf(FilteringMessageListenerAdapter.class);
|
||||
|
||||
Object delegate = TestUtils.getPropertyValue(messageListener, "delegate");
|
||||
|
||||
@@ -123,7 +123,7 @@ public class KafkaMessageDrivenChannelAdapterParserTests {
|
||||
adapter.afterPropertiesSet();
|
||||
|
||||
messageListener = containerProps.getMessageListener();
|
||||
assertThat(messageListener).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class);
|
||||
assertThat(messageListener).isInstanceOf(RetryingMessageListenerAdapter.class);
|
||||
|
||||
delegate = TestUtils.getPropertyValue(messageListener, "delegate");
|
||||
|
||||
@@ -133,21 +133,21 @@ public class KafkaMessageDrivenChannelAdapterParserTests {
|
||||
adapter.afterPropertiesSet();
|
||||
|
||||
messageListener = containerProps.getMessageListener();
|
||||
assertThat(messageListener).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class);
|
||||
assertThat(messageListener).isInstanceOf(FilteringMessageListenerAdapter.class);
|
||||
|
||||
delegate = TestUtils.getPropertyValue(messageListener, "delegate");
|
||||
|
||||
assertThat(delegate).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class);
|
||||
assertThat(delegate).isInstanceOf(RetryingMessageListenerAdapter.class);
|
||||
|
||||
adapter.setFilterInRetry(true);
|
||||
adapter.afterPropertiesSet();
|
||||
|
||||
messageListener = containerProps.getMessageListener();
|
||||
assertThat(messageListener).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class);
|
||||
assertThat(messageListener).isInstanceOf(RetryingMessageListenerAdapter.class);
|
||||
|
||||
delegate = TestUtils.getPropertyValue(messageListener, "delegate");
|
||||
|
||||
assertThat(delegate).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class);
|
||||
assertThat(delegate).isInstanceOf(FilteringMessageListenerAdapter.class);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -23,6 +23,7 @@ import java.util.Arrays;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import org.apache.kafka.clients.consumer.Consumer;
|
||||
import org.apache.kafka.clients.consumer.ConsumerConfig;
|
||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.apache.kafka.clients.producer.ProducerRecord;
|
||||
@@ -89,8 +90,9 @@ public class MessageDrivenAdapterTests {
|
||||
adapter.setRecordMessageConverter(new MessagingMessageConverter() {
|
||||
|
||||
@Override
|
||||
public Message<?> toMessage(ConsumerRecord<?, ?> record, Acknowledgment acknowledgment, Type type) {
|
||||
Message<?> message = super.toMessage(record, acknowledgment, type);
|
||||
public Message<?> toMessage(ConsumerRecord<?, ?> record, Acknowledgment acknowledgment,
|
||||
Consumer<?, ?> consumer, Type type) {
|
||||
Message<?> message = super.toMessage(record, acknowledgment, consumer, type);
|
||||
return MessageBuilder.fromMessage(message).setHeader("testHeader", "testValue").build();
|
||||
}
|
||||
|
||||
@@ -136,7 +138,8 @@ public class MessageDrivenAdapterTests {
|
||||
adapter.setMessageConverter(new RecordMessageConverter() {
|
||||
|
||||
@Override
|
||||
public Message<?> toMessage(ConsumerRecord<?, ?> record, Acknowledgment acknowledgment, Type payloadType) {
|
||||
public Message<?> toMessage(ConsumerRecord<?, ?> record, Acknowledgment acknowledgment,
|
||||
Consumer<?, ?> consumer, Type type) {
|
||||
throw new RuntimeException("testError");
|
||||
}
|
||||
|
||||
@@ -173,8 +176,9 @@ public class MessageDrivenAdapterTests {
|
||||
adapter.setBatchMessageConverter(new BatchMessagingMessageConverter() {
|
||||
|
||||
@Override
|
||||
public Message<?> toMessage(List<ConsumerRecord<?, ?>> records, Acknowledgment acknowledgment, Type type) {
|
||||
Message<?> message = super.toMessage(records, acknowledgment, type);
|
||||
public Message<?> toMessage(List<ConsumerRecord<?, ?>> records, Acknowledgment acknowledgment,
|
||||
Consumer<?, ?> consumer, Type type) {
|
||||
Message<?> message = super.toMessage(records, acknowledgment, consumer, type);
|
||||
return MessageBuilder.fromMessage(message).setHeader("testHeader", "testValue").build();
|
||||
}
|
||||
|
||||
@@ -211,7 +215,7 @@ public class MessageDrivenAdapterTests {
|
||||
|
||||
@Override
|
||||
public Message<?> toMessage(List<ConsumerRecord<?, ?>> records, Acknowledgment acknowledgment,
|
||||
Type payloadType) {
|
||||
Consumer<?, ?> consumer, Type payloadType) {
|
||||
throw new RuntimeException("testError");
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user