Add retry and DLQ support for the Kafka binder
Fixes #498 - Honour the retry settings from ConsumerProperties; - Add an additional `enableDlq` option in KafkaConsumerProperties that enables forwarding failed messages to a DLQ topic.
This commit is contained in:
committed by
Gary Russell
parent
4247fe42bd
commit
fa42b62074
@@ -31,6 +31,8 @@ public class KafkaConsumerProperties {
|
||||
|
||||
private KafkaMessageChannelBinder.StartOffset startOffset = null;
|
||||
|
||||
private boolean enableDlq = false;
|
||||
|
||||
public void setMinPartitionCount(int minPartitionCount) {
|
||||
this.minPartitionCount = minPartitionCount;
|
||||
}
|
||||
@@ -63,4 +65,12 @@ public class KafkaConsumerProperties {
|
||||
public void setStartOffset(KafkaMessageChannelBinder.StartOffset startOffset) {
|
||||
this.startOffset = startOffset;
|
||||
}
|
||||
|
||||
public boolean isEnableDlq() {
|
||||
return enableDlq;
|
||||
}
|
||||
|
||||
public void setEnableDlq(boolean enableDlq) {
|
||||
this.enableDlq = enableDlq;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,15 +32,14 @@ import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ThreadFactory;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import kafka.admin.AdminUtils;
|
||||
import kafka.api.OffsetRequest;
|
||||
import kafka.serializer.Decoder;
|
||||
import kafka.serializer.DefaultDecoder;
|
||||
import kafka.utils.ZKStringSerializer$;
|
||||
import kafka.utils.ZkUtils;
|
||||
import org.I0Itec.zkclient.ZkClient;
|
||||
import org.apache.kafka.clients.producer.Callback;
|
||||
import org.apache.kafka.clients.producer.Producer;
|
||||
import org.apache.kafka.clients.producer.ProducerConfig;
|
||||
import org.apache.kafka.clients.producer.ProducerRecord;
|
||||
import org.apache.kafka.clients.producer.RecordMetadata;
|
||||
import org.apache.kafka.common.serialization.ByteArraySerializer;
|
||||
import org.apache.kafka.common.utils.Utils;
|
||||
|
||||
import org.springframework.beans.factory.DisposableBean;
|
||||
import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
|
||||
@@ -63,12 +62,16 @@ import org.springframework.integration.handler.AbstractMessageHandler;
|
||||
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
|
||||
import org.springframework.integration.kafka.core.ConnectionFactory;
|
||||
import org.springframework.integration.kafka.core.DefaultConnectionFactory;
|
||||
import org.springframework.integration.kafka.core.KafkaMessage;
|
||||
import org.springframework.integration.kafka.core.Partition;
|
||||
import org.springframework.integration.kafka.core.ZookeeperConfiguration;
|
||||
import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter;
|
||||
import org.springframework.integration.kafka.listener.AcknowledgingMessageListener;
|
||||
import org.springframework.integration.kafka.listener.Acknowledgment;
|
||||
import org.springframework.integration.kafka.listener.ErrorHandler;
|
||||
import org.springframework.integration.kafka.listener.KafkaMessageListenerContainer;
|
||||
import org.springframework.integration.kafka.listener.KafkaNativeOffsetManager;
|
||||
import org.springframework.integration.kafka.listener.MessageListener;
|
||||
import org.springframework.integration.kafka.listener.OffsetManager;
|
||||
import org.springframework.integration.kafka.support.KafkaHeaders;
|
||||
import org.springframework.integration.kafka.support.ProducerConfiguration;
|
||||
@@ -90,8 +93,15 @@ import org.springframework.retry.policy.SimpleRetryPolicy;
|
||||
import org.springframework.retry.support.RetryTemplate;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
import org.springframework.util.ObjectUtils;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
import kafka.admin.AdminUtils;
|
||||
import kafka.api.OffsetRequest;
|
||||
import kafka.serializer.Decoder;
|
||||
import kafka.serializer.DefaultDecoder;
|
||||
import kafka.utils.ZKStringSerializer$;
|
||||
import kafka.utils.ZkUtils;
|
||||
import scala.collection.Seq;
|
||||
|
||||
/**
|
||||
@@ -105,7 +115,12 @@ import scala.collection.Seq;
|
||||
* @author Mark Fisher
|
||||
* @author Soby Chacko
|
||||
*/
|
||||
public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, ExtendedConsumerProperties<KafkaConsumerProperties>, ExtendedProducerProperties<KafkaProducerProperties>> implements ExtendedPropertiesBinder<MessageChannel, KafkaConsumerProperties, KafkaProducerProperties> {
|
||||
public class KafkaMessageChannelBinder
|
||||
extends
|
||||
AbstractBinder<MessageChannel, ExtendedConsumerProperties<KafkaConsumerProperties>,
|
||||
ExtendedProducerProperties<KafkaProducerProperties>>
|
||||
implements ExtendedPropertiesBinder<MessageChannel, KafkaConsumerProperties, KafkaProducerProperties>,
|
||||
DisposableBean {
|
||||
|
||||
public static final ByteArraySerializer BYTE_ARRAY_SERIALIZER = new ByteArraySerializer();
|
||||
|
||||
@@ -153,19 +168,20 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
|
||||
private ProducerListener producerListener;
|
||||
|
||||
private volatile Producer<byte[],byte[]> dlqProducer;
|
||||
|
||||
private KafkaExtendedBindingProperties extendedBindingProperties = new KafkaExtendedBindingProperties();
|
||||
|
||||
public KafkaMessageChannelBinder(ZookeeperConnect zookeeperConnect, String brokers,
|
||||
String zkAddress, String... headersToMap) {
|
||||
public KafkaMessageChannelBinder(ZookeeperConnect zookeeperConnect, String brokers, String zkAddress,
|
||||
String... headersToMap) {
|
||||
this.zookeeperConnect = zookeeperConnect;
|
||||
this.brokers = brokers;
|
||||
this.zkAddress = zkAddress;
|
||||
if (headersToMap.length > 0) {
|
||||
String[] combinedHeadersToMap = Arrays.copyOfRange(
|
||||
BinderHeaders.STANDARD_HEADERS, 0,
|
||||
String[] combinedHeadersToMap = Arrays.copyOfRange(BinderHeaders.STANDARD_HEADERS, 0,
|
||||
BinderHeaders.STANDARD_HEADERS.length + headersToMap.length);
|
||||
System.arraycopy(headersToMap, 0, combinedHeadersToMap,
|
||||
BinderHeaders.STANDARD_HEADERS.length, headersToMap.length);
|
||||
System.arraycopy(headersToMap, 0, combinedHeadersToMap, BinderHeaders.STANDARD_HEADERS.length,
|
||||
headersToMap.length);
|
||||
this.headersToMap = combinedHeadersToMap;
|
||||
}
|
||||
else {
|
||||
@@ -218,8 +234,7 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
ZookeeperConfiguration configuration = new ZookeeperConfiguration(this.zookeeperConnect);
|
||||
configuration.setBufferSize(socketBufferSize);
|
||||
configuration.setMaxWait(maxWait);
|
||||
DefaultConnectionFactory defaultConnectionFactory =
|
||||
new DefaultConnectionFactory(configuration);
|
||||
DefaultConnectionFactory defaultConnectionFactory = new DefaultConnectionFactory(configuration);
|
||||
defaultConnectionFactory.afterPropertiesSet();
|
||||
this.connectionFactory = defaultConnectionFactory;
|
||||
if (retryOperations == null) {
|
||||
@@ -231,13 +246,21 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
|
||||
ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy();
|
||||
backOffPolicy.setInitialInterval(100);
|
||||
backOffPolicy.setMultiplier((double) 2);
|
||||
backOffPolicy.setMultiplier(2);
|
||||
backOffPolicy.setMaxInterval(1000);
|
||||
retryTemplate.setBackOffPolicy(backOffPolicy);
|
||||
retryOperations = retryTemplate;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void destroy() throws Exception {
|
||||
if (dlqProducer != null) {
|
||||
dlqProducer.close();
|
||||
dlqProducer = null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Allowed chars are ASCII alphanumerics, '.', '_' and '-'.
|
||||
*/
|
||||
@@ -247,7 +270,8 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
for (byte b : utf8) {
|
||||
if (!((b >= 'a') && (b <= 'z') || (b >= 'A') && (b <= 'Z') || (b >= '0') && (b <= '9') || (b == '.')
|
||||
|| (b == '-') || (b == '_'))) {
|
||||
throw new IllegalArgumentException("Topic name can only have ASCII alphanumerics, '.', '_' and '-'");
|
||||
throw new IllegalArgumentException(
|
||||
"Topic name can only have ASCII alphanumerics, '.', '_' and '-'");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -313,22 +337,29 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
@Override
|
||||
protected Binding<MessageChannel> doBindConsumer(String name, String group, MessageChannel inputChannel,
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> properties) {
|
||||
// If the caller provides a consumer group, use it; otherwise an anonymous consumer group
|
||||
// is generated each time, such that each anonymous binding will receive all messages.
|
||||
// Consumers reset offsets at the latest time by default, which allows them to receive only
|
||||
// If the caller provides a consumer group, use it; otherwise an anonymous
|
||||
// consumer group
|
||||
// is generated each time, such that each anonymous binding will receive all
|
||||
// messages.
|
||||
// Consumers reset offsets at the latest time by default, which allows them to
|
||||
// receive only
|
||||
// messages sent after they've been bound. That behavior can be changed with the
|
||||
// "resetOffsets" and "startOffset" properties.
|
||||
boolean anonymous = !StringUtils.hasText(group);
|
||||
Assert.isTrue(!anonymous || !properties.getExtension().isEnableDlq(),
|
||||
"DLQ support is not available for anonymous subscriptions");
|
||||
String consumerGroup = anonymous ? "anonymous." + UUID.randomUUID().toString() : group;
|
||||
// The reference point, if not set explicitly is the latest time for anonymous subscriptions and the
|
||||
// earliest time for group subscriptions. This allows the latter to receive messages published before the group
|
||||
// The reference point, if not set explicitly is the latest time for anonymous
|
||||
// subscriptions and the
|
||||
// earliest time for group subscriptions. This allows the latter to receive
|
||||
// messages published before the group
|
||||
// has been created.
|
||||
long referencePoint = properties.getExtension().getStartOffset() != null ?
|
||||
properties.getExtension().getStartOffset().getReferencePoint() : (anonymous ? OffsetRequest.LatestTime() : OffsetRequest.EarliestTime());
|
||||
long referencePoint = properties.getExtension().getStartOffset() != null
|
||||
? properties.getExtension().getStartOffset().getReferencePoint()
|
||||
: (anonymous ? OffsetRequest.LatestTime() : OffsetRequest.EarliestTime());
|
||||
return createKafkaConsumer(name, inputChannel, properties, consumerGroup, referencePoint);
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public Binding<MessageChannel> doBindProducer(String name, MessageChannel moduleOutputChannel,
|
||||
ExtendedProducerProperties<KafkaProducerProperties> properties) {
|
||||
@@ -352,8 +383,10 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
producerMetadata.setBatchBytes(properties.getExtension().getBufferSize());
|
||||
Properties additionalProps = new Properties();
|
||||
additionalProps.put(ProducerConfig.ACKS_CONFIG, String.valueOf(requiredAcks));
|
||||
additionalProps.put(ProducerConfig.LINGER_MS_CONFIG, String.valueOf(properties.getExtension().getBatchTimeout()));
|
||||
ProducerFactoryBean<byte[], byte[]> producerFB = new ProducerFactoryBean<>(producerMetadata, brokers, additionalProps);
|
||||
additionalProps.put(ProducerConfig.LINGER_MS_CONFIG,
|
||||
String.valueOf(properties.getExtension().getBatchTimeout()));
|
||||
ProducerFactoryBean<byte[], byte[]> producerFB = new ProducerFactoryBean<>(producerMetadata, brokers,
|
||||
additionalProps);
|
||||
|
||||
try {
|
||||
final ProducerConfiguration<byte[], byte[]> producerConfiguration = new ProducerConfiguration<>(
|
||||
@@ -361,12 +394,19 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
producerConfiguration.setProducerListener(producerListener);
|
||||
|
||||
MessageHandler handler = new SendingHandler(name, properties, partitions.size(), producerConfiguration);
|
||||
EventDrivenConsumer consumer = new EventDrivenConsumer((SubscribableChannel) moduleOutputChannel,
|
||||
handler);
|
||||
EventDrivenConsumer consumer = new EventDrivenConsumer((SubscribableChannel) moduleOutputChannel, handler) {
|
||||
|
||||
@Override
|
||||
protected void doStop() {
|
||||
super.doStop();
|
||||
producerConfiguration.stop();
|
||||
}
|
||||
};
|
||||
consumer.setBeanFactory(this.getBeanFactory());
|
||||
consumer.setBeanName("outbound." + name);
|
||||
consumer.afterPropertiesSet();
|
||||
DefaultBinding<MessageChannel> producerBinding = new DefaultBinding<>(name, null, moduleOutputChannel, consumer);
|
||||
DefaultBinding<MessageChannel> producerBinding = new DefaultBinding<>(name, null, moduleOutputChannel,
|
||||
consumer);
|
||||
consumer.start();
|
||||
return producerBinding;
|
||||
}
|
||||
@@ -376,7 +416,8 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates a Kafka topic if needed, or try to increase its partition count to the desired number.
|
||||
* Creates a Kafka topic if needed, or try to increase its partition count to the
|
||||
* desired number.
|
||||
*/
|
||||
private Collection<Partition> ensureTopicCreated(final String topicName, final int numPartitions,
|
||||
int replicationFactor) {
|
||||
@@ -388,9 +429,8 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
// createOrUpdateTopicPartitionAssignmentPathInZK(..., update=true)
|
||||
final Properties topicConfig = new Properties();
|
||||
Seq<Object> brokerList = ZkUtils.getSortedBrokerList(zkClient);
|
||||
final scala.collection.Map<Object, Seq<Object>> replicaAssignment = AdminUtils.assignReplicasToBrokers
|
||||
(brokerList,
|
||||
numPartitions, replicationFactor, -1, -1);
|
||||
final scala.collection.Map<Object, Seq<Object>> replicaAssignment = AdminUtils
|
||||
.assignReplicasToBrokers(brokerList, numPartitions, replicationFactor, -1, -1);
|
||||
retryOperations.execute(new RetryCallback<Object, RuntimeException>() {
|
||||
|
||||
@Override
|
||||
@@ -401,21 +441,22 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
}
|
||||
});
|
||||
try {
|
||||
Collection<Partition> partitions = retryOperations.execute(new RetryCallback<Collection<Partition>, Exception>() {
|
||||
Collection<Partition> partitions = retryOperations
|
||||
.execute(new RetryCallback<Collection<Partition>, Exception>() {
|
||||
|
||||
@Override
|
||||
public Collection<Partition> doWithRetry(RetryContext context) throws Exception {
|
||||
connectionFactory.refreshMetadata(Collections.singleton(topicName));
|
||||
Collection<Partition> partitions = connectionFactory.getPartitions(topicName);
|
||||
if (partitions.size() < numPartitions) {
|
||||
throw new IllegalStateException("The number of expected partitions was: " + numPartitions
|
||||
+ ", but " +
|
||||
partitions.size() + " have been found instead");
|
||||
}
|
||||
connectionFactory.getLeaders(partitions);
|
||||
return partitions;
|
||||
}
|
||||
});
|
||||
@Override
|
||||
public Collection<Partition> doWithRetry(RetryContext context) throws Exception {
|
||||
connectionFactory.refreshMetadata(Collections.singleton(topicName));
|
||||
Collection<Partition> partitions = connectionFactory.getPartitions(topicName);
|
||||
if (partitions.size() < numPartitions) {
|
||||
throw new IllegalStateException(
|
||||
"The number of expected partitions was: " + numPartitions + ", but "
|
||||
+ partitions.size() + " have been found instead");
|
||||
}
|
||||
connectionFactory.getLeaders(partitions);
|
||||
return partitions;
|
||||
}
|
||||
});
|
||||
return partitions;
|
||||
}
|
||||
catch (Exception e) {
|
||||
@@ -438,7 +479,7 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
if (instance == 0) {
|
||||
throw new IllegalArgumentException("Instance count cannot be zero");
|
||||
}
|
||||
int numPartitions = Math.max(minKafkaPartitions, instance * properties.getConcurrency());
|
||||
final int numPartitions = Math.max(minKafkaPartitions, instance * properties.getConcurrency());
|
||||
Collection<Partition> allPartitions = ensureTopicCreated(name, numPartitions, replicationFactor);
|
||||
|
||||
Decoder<byte[]> valueDecoder = new DefaultDecoder(null);
|
||||
@@ -466,8 +507,8 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
bridge.setBeanName("bridge." + name);
|
||||
|
||||
Assert.isTrue(!CollectionUtils.isEmpty(listenedPartitions), "A list of partitions must be provided");
|
||||
final KafkaMessageListenerContainer messageListenerContainer = new KafkaMessageListenerContainer(connectionFactory,
|
||||
listenedPartitions.toArray(new Partition[listenedPartitions.size()]));
|
||||
final KafkaMessageListenerContainer messageListenerContainer = new KafkaMessageListenerContainer(
|
||||
connectionFactory, listenedPartitions.toArray(new Partition[listenedPartitions.size()]));
|
||||
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Listened partitions: " + StringUtils.collectionToCommaDelimitedString(listenedPartitions));
|
||||
@@ -485,17 +526,97 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
messageListenerContainer.setConcurrency(concurrency);
|
||||
final ExecutorService dispatcherTaskExecutor = Executors.newFixedThreadPool(concurrency, DAEMON_THREAD_FACTORY);
|
||||
messageListenerContainer.setDispatcherTaskExecutor(dispatcherTaskExecutor);
|
||||
|
||||
final KafkaMessageDrivenChannelAdapter kafkaMessageDrivenChannelAdapter =
|
||||
new KafkaMessageDrivenChannelAdapter(messageListenerContainer);
|
||||
|
||||
final KafkaMessageDrivenChannelAdapter kafkaMessageDrivenChannelAdapter = new KafkaMessageDrivenChannelAdapter(
|
||||
messageListenerContainer);
|
||||
kafkaMessageDrivenChannelAdapter.setBeanFactory(this.getBeanFactory());
|
||||
kafkaMessageDrivenChannelAdapter.setKeyDecoder(keyDecoder);
|
||||
kafkaMessageDrivenChannelAdapter.setPayloadDecoder(valueDecoder);
|
||||
kafkaMessageDrivenChannelAdapter.setOutputChannel(bridge);
|
||||
kafkaMessageDrivenChannelAdapter.setAutoCommitOffset(properties.getExtension().isAutoCommitOffset());
|
||||
kafkaMessageDrivenChannelAdapter.afterPropertiesSet();
|
||||
kafkaMessageDrivenChannelAdapter.start();
|
||||
|
||||
// we need to wrap the adapter listener into a retrying listener so that the retry
|
||||
// logic is applied before the ErrorHandler is executed
|
||||
final RetryTemplate retryTemplate = buildRetryTemplateIfRetryEnabled(properties);
|
||||
if (retryTemplate != null) {
|
||||
if (properties.getExtension().isAutoCommitOffset()) {
|
||||
final MessageListener originalMessageListener = (MessageListener) messageListenerContainer
|
||||
.getMessageListener();
|
||||
messageListenerContainer.setMessageListener(new MessageListener() {
|
||||
@Override
|
||||
public void onMessage(final KafkaMessage message) {
|
||||
try {
|
||||
retryTemplate.execute(new RetryCallback<Object, Throwable>() {
|
||||
@Override
|
||||
public Object doWithRetry(RetryContext context) {
|
||||
originalMessageListener.onMessage(message);
|
||||
return null;
|
||||
}
|
||||
});
|
||||
}
|
||||
catch (Throwable throwable) {
|
||||
if (throwable instanceof RuntimeException) {
|
||||
throw (RuntimeException) throwable;
|
||||
}
|
||||
else {
|
||||
throw new RuntimeException(throwable);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
else {
|
||||
messageListenerContainer.setMessageListener(new AcknowledgingMessageListener() {
|
||||
final AcknowledgingMessageListener originalMessageListener = (AcknowledgingMessageListener) messageListenerContainer
|
||||
.getMessageListener();
|
||||
@Override
|
||||
public void onMessage(final KafkaMessage message, final Acknowledgment acknowledgment) {
|
||||
retryTemplate.execute(new RetryCallback<Object, RuntimeException>() {
|
||||
@Override
|
||||
public Object doWithRetry(RetryContext context) {
|
||||
originalMessageListener.onMessage(message, acknowledgment);
|
||||
return null;
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
if (properties.getExtension().isEnableDlq()) {
|
||||
final String dlqTopic = "error." + name + "." + group;
|
||||
initDlqProducer();
|
||||
messageListenerContainer.setErrorHandler(new ErrorHandler() {
|
||||
@Override
|
||||
public void handle(Exception thrownException, final KafkaMessage message) {
|
||||
final byte[] key = message.getMessage().key() != null ? Utils.toArray(message.getMessage().key())
|
||||
: null;
|
||||
final byte[] payload = message.getMessage().payload() != null
|
||||
? Utils.toArray(message.getMessage().payload()) : null;
|
||||
dlqProducer.send(new ProducerRecord<>(dlqTopic, key, payload), new Callback() {
|
||||
@Override
|
||||
public void onCompletion(RecordMetadata metadata, Exception exception) {
|
||||
StringBuffer messageLog = new StringBuffer();
|
||||
messageLog.append(" a message with key='"
|
||||
+ toDisplayString(ObjectUtils.nullSafeToString(key), 50) + "'");
|
||||
messageLog.append(" and payload='"
|
||||
+ toDisplayString(ObjectUtils.nullSafeToString(payload), 50) + "'");
|
||||
messageLog.append(" received from " + message.getMetadata().getPartition());
|
||||
if (exception != null) {
|
||||
logger.error("Error sending to DLQ" + messageLog.toString(), exception);
|
||||
}
|
||||
else {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Sent to DLQ " + messageLog.toString());
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
}
|
||||
kafkaMessageDrivenChannelAdapter.start();
|
||||
|
||||
EventDrivenConsumer edc = new EventDrivenConsumer(bridge, rh) {
|
||||
|
||||
@@ -529,12 +650,39 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
return consumerBinding;
|
||||
}
|
||||
|
||||
private synchronized void initDlqProducer() {
|
||||
try {
|
||||
if (dlqProducer == null) {
|
||||
synchronized (this) {
|
||||
if (dlqProducer == null) {
|
||||
// we can use the producer defaults as we do not need to tune
|
||||
// performance
|
||||
ProducerMetadata<byte[], byte[]> producerMetadata = new ProducerMetadata<>("dlqKafkaProducer",
|
||||
byte[].class, byte[].class, BYTE_ARRAY_SERIALIZER, BYTE_ARRAY_SERIALIZER);
|
||||
producerMetadata.setSync(false);
|
||||
producerMetadata.setCompressionType(ProducerMetadata.CompressionType.none);
|
||||
producerMetadata.setBatchBytes(16384);
|
||||
Properties additionalProps = new Properties();
|
||||
additionalProps.put(ProducerConfig.ACKS_CONFIG, String.valueOf(requiredAcks));
|
||||
additionalProps.put(ProducerConfig.LINGER_MS_CONFIG,
|
||||
String.valueOf(0));
|
||||
ProducerFactoryBean<byte[], byte[]> producerFactoryBean = new ProducerFactoryBean<>(
|
||||
producerMetadata, brokers, additionalProps);
|
||||
dlqProducer = producerFactoryBean.getObject();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new RuntimeException("Cannot initialize DLQ producer:", e);
|
||||
}
|
||||
}
|
||||
|
||||
private OffsetManager createOffsetManager(String group, long referencePoint) {
|
||||
try {
|
||||
|
||||
KafkaNativeOffsetManager kafkaOffsetManager =
|
||||
new KafkaNativeOffsetManager(connectionFactory, zookeeperConnect,
|
||||
Collections.<Partition, Long>emptyMap());
|
||||
KafkaNativeOffsetManager kafkaOffsetManager = new KafkaNativeOffsetManager(connectionFactory,
|
||||
zookeeperConnect, Collections.<Partition, Long>emptyMap());
|
||||
kafkaOffsetManager.setConsumerId(group);
|
||||
kafkaOffsetManager.setReferenceTimestamp(referencePoint);
|
||||
kafkaOffsetManager.afterPropertiesSet();
|
||||
@@ -552,14 +700,22 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private String toDisplayString(String original, int maxCharacters) {
|
||||
if (original.length() <= maxCharacters) {
|
||||
return original;
|
||||
}
|
||||
return original.substring(0, maxCharacters) + "...";
|
||||
}
|
||||
|
||||
@Override
|
||||
public void doManualAck(LinkedList<MessageHeaders> messageHeadersList) {
|
||||
Iterator<MessageHeaders> iterator = messageHeadersList.iterator();
|
||||
while (iterator.hasNext()) {
|
||||
MessageHeaders headers = iterator.next();
|
||||
Acknowledgment acknowledgment = (Acknowledgment) headers.get(KafkaHeaders.ACKNOWLEDGMENT);
|
||||
Assert.notNull(acknowledgment, "Acknowledgement shouldn't be null when acknowledging kafka message " +
|
||||
"manually.");
|
||||
Assert.notNull(acknowledgment,
|
||||
"Acknowledgement shouldn't be null when acknowledging kafka message " + "manually.");
|
||||
acknowledgment.acknowledge();
|
||||
}
|
||||
}
|
||||
@@ -576,7 +732,7 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
|
||||
private class ReceivingHandler extends AbstractReplyProducingMessageHandler {
|
||||
|
||||
private ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties;
|
||||
private final ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties;
|
||||
|
||||
public ReceivingHandler(ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties) {
|
||||
this.consumerProperties = consumerProperties;
|
||||
@@ -587,8 +743,7 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
protected Object handleRequestMessage(Message<?> requestMessage) {
|
||||
if (HeaderMode.embeddedHeaders.equals(consumerProperties.getHeaderMode())) {
|
||||
MessageValues messageValues = extractMessageValues(requestMessage);
|
||||
return MessageBuilder.createMessage(messageValues.getPayload(), new KafkaBinderHeaders(
|
||||
messageValues));
|
||||
return MessageBuilder.createMessage(messageValues.getPayload(), new KafkaBinderHeaders(messageValues));
|
||||
}
|
||||
else {
|
||||
return requestMessage;
|
||||
@@ -603,7 +758,6 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
protected boolean shouldCopyRequestHeaders() {
|
||||
// prevent the message from being copied again in superclass
|
||||
@@ -626,16 +780,14 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
private final PartitionHandler partitionHandler;
|
||||
|
||||
private SendingHandler(String topicName, ExtendedProducerProperties<KafkaProducerProperties> properties,
|
||||
int numberOfPartitions,
|
||||
ProducerConfiguration<byte[], byte[]> producerConfiguration) {
|
||||
int numberOfPartitions, ProducerConfiguration<byte[], byte[]> producerConfiguration) {
|
||||
this.topicName = topicName;
|
||||
producerProperties = properties;
|
||||
this.numberOfKafkaPartitions = numberOfPartitions;
|
||||
ConfigurableListableBeanFactory beanFactory = KafkaMessageChannelBinder.this.getBeanFactory();
|
||||
this.setBeanFactory(beanFactory);
|
||||
this.producerConfiguration = producerConfiguration;
|
||||
this.partitionHandler = new PartitionHandler(beanFactory, evaluationContext, partitionSelector,
|
||||
properties);
|
||||
this.partitionHandler = new PartitionHandler(beanFactory, evaluationContext, partitionSelector, properties);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -655,8 +807,7 @@ public class KafkaMessageChannelBinder extends AbstractBinder<MessageChannel, Ex
|
||||
}
|
||||
else if (HeaderMode.raw.equals(producerProperties.getHeaderMode())) {
|
||||
Object contentType = message.getHeaders().get(MessageHeaders.CONTENT_TYPE);
|
||||
if (contentType != null
|
||||
&& !contentType.equals(MediaType.APPLICATION_OCTET_STREAM_VALUE)) {
|
||||
if (contentType != null && !contentType.equals(MediaType.APPLICATION_OCTET_STREAM_VALUE)) {
|
||||
logger.error("Raw mode supports only " + MediaType.APPLICATION_OCTET_STREAM_VALUE + " content type"
|
||||
+ message.getPayload().getClass());
|
||||
}
|
||||
|
||||
@@ -21,6 +21,7 @@ import static org.hamcrest.Matchers.not;
|
||||
import static org.hamcrest.Matchers.nullValue;
|
||||
import static org.hamcrest.collection.IsCollectionWithSize.hasSize;
|
||||
import static org.junit.Assert.assertArrayEquals;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertThat;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
@@ -50,7 +51,9 @@ import org.springframework.integration.kafka.support.ZookeeperConnect;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
import org.springframework.messaging.MessageHandler;
|
||||
import org.springframework.messaging.MessagingException;
|
||||
import org.springframework.messaging.support.GenericMessage;
|
||||
import org.springframework.messaging.support.MessageBuilder;
|
||||
|
||||
|
||||
/**
|
||||
@@ -115,6 +118,48 @@ public class KafkaBinderTests extends PartitionCapableBinderTests<KafkaTestBinde
|
||||
throw new UnsupportedOperationException("'spyOn' is not used by Kafka tests");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testDlqAndRetry() {
|
||||
KafkaTestBinder binder = getBinder();
|
||||
|
||||
DirectChannel moduleOutputChannel = new DirectChannel();
|
||||
DirectChannel moduleInputChannel = new DirectChannel();
|
||||
QueueChannel dlqChannel = new QueueChannel();
|
||||
FailingInvocationCountingMessageHandler handler = new FailingInvocationCountingMessageHandler();
|
||||
moduleInputChannel.subscribe(handler);
|
||||
ExtendedProducerProperties<KafkaProducerProperties> producerProperties = createProducerProperties();
|
||||
producerProperties.setPartitionCount(10);
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> consumerProperties = createConsumerProperties();
|
||||
consumerProperties.setMaxAttempts(3);
|
||||
consumerProperties.setBackOffInitialInterval(100);
|
||||
consumerProperties.setBackOffMaxInterval(150);
|
||||
consumerProperties.getExtension().setMinPartitionCount(10);
|
||||
consumerProperties.getExtension().setEnableDlq(true);
|
||||
long uniqueBindingId = System.currentTimeMillis();
|
||||
Binding<MessageChannel> producerBinding = binder.bindProducer("retryTest." + uniqueBindingId + ".0",
|
||||
moduleOutputChannel, producerProperties);
|
||||
Binding<MessageChannel> consumerBinding = binder.bindConsumer("retryTest." + uniqueBindingId + ".0", "testGroup",
|
||||
moduleInputChannel, consumerProperties);
|
||||
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> dlqConsumerProperties = createConsumerProperties();
|
||||
dlqConsumerProperties.setMaxAttempts(1);
|
||||
|
||||
Binding<MessageChannel> dlqConsumerBinding = binder.bindConsumer(
|
||||
"error.retryTest." + uniqueBindingId + ".0.testGroup", null, dlqChannel, dlqConsumerProperties);
|
||||
|
||||
String testMessagePayload = "test." + UUID.randomUUID().toString();
|
||||
Message<String> testMessage = MessageBuilder.withPayload(testMessagePayload).build();
|
||||
moduleOutputChannel.send(testMessage);
|
||||
|
||||
Message<?> receivedMessage = receive(dlqChannel, 3);
|
||||
assertNotNull(receivedMessage);
|
||||
assertEquals(testMessagePayload, receivedMessage.getPayload());
|
||||
assertThat(handler.getInvocationCount(), equalTo(consumerProperties.getMaxAttempts()));
|
||||
dlqConsumerBinding.unbind();
|
||||
consumerBinding.unbind();
|
||||
producerBinding.unbind();
|
||||
}
|
||||
|
||||
@Test(expected = IllegalArgumentException.class)
|
||||
public void testValidateKafkaTopicName() {
|
||||
KafkaMessageChannelBinder.validateTopicName("foo:bar");
|
||||
@@ -400,4 +445,22 @@ public class KafkaBinderTests extends PartitionCapableBinderTests<KafkaTestBinde
|
||||
assertTrue("Kafka Sync Producer should have been enabled.", producerConfiguration.getProducerMetadata().isSync());
|
||||
producerBinding.unbind();
|
||||
}
|
||||
|
||||
private static class FailingInvocationCountingMessageHandler implements MessageHandler {
|
||||
|
||||
private int invocationCount = 0;
|
||||
|
||||
public FailingInvocationCountingMessageHandler() {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void handleMessage(Message<?> message) throws MessagingException {
|
||||
invocationCount++;
|
||||
throw new RuntimeException();
|
||||
}
|
||||
|
||||
public int getInvocationCount() {
|
||||
return invocationCount;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -63,7 +63,17 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
|
||||
* waiting up to 1s (times the {@link #timeoutMultiplier}).
|
||||
*/
|
||||
protected Message<?> receive(PollableChannel channel) {
|
||||
return channel.receive((int)(1000 * timeoutMultiplier));
|
||||
return receive(channel, 1);
|
||||
}
|
||||
|
||||
/**
|
||||
* Attempt to receive a message on the given channel,
|
||||
* waiting up to 1s * additionalMultiplier * {@link #timeoutMultiplier}).
|
||||
*
|
||||
* Allows accomodating tests which are slower than normal (e.g. retry).
|
||||
*/
|
||||
protected Message<?> receive(PollableChannel channel, int additionalMultiplier) {
|
||||
return channel.receive((int)(1000 * timeoutMultiplier * additionalMultiplier));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user