From 99db7dff9b01cc7533c36b51a009f49bb3a1d61b Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Tue, 31 May 2016 17:19:54 -0400 Subject: [PATCH] GH-541 Make reconnect time for Kafka/Rabbit configurable --- .../spring-cloud-stream-binder-kafka/pom.xml | 2 +- .../stream/binder/kafka/KafkaConsumerProperties.java | 10 ++++++++++ .../stream/binder/kafka/KafkaMessageChannelBinder.java | 1 + .../stream/binder/rabbit/RabbitConsumerProperties.java | 10 ++++++++++ .../binder/rabbit/RabbitMessageChannelBinder.java | 5 ++--- .../main/asciidoc/spring-cloud-stream-overview.adoc | 8 ++++++++ 6 files changed, 32 insertions(+), 4 deletions(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/pom.xml b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/pom.xml index 5c11fda2e..e2d5cbbfa 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/pom.xml +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/pom.xml @@ -16,7 +16,7 @@ 0.8.2.2 2.6.0 - 1.3.0.RELEASE + 1.3.1.BUILD-SNAPSHOT 1.0.0 diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java index 0f3ff40c2..c52a38221 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java @@ -31,6 +31,8 @@ public class KafkaConsumerProperties { private boolean enableDlq; + private int recoveryInterval = 5000; + public boolean isAutoCommitOffset() { return autoCommitOffset; } @@ -70,4 +72,12 @@ public class KafkaConsumerProperties { public void setAutoCommitOnError(Boolean autoCommitOnError) { this.autoCommitOnError = autoCommitOnError; } + + public int getRecoveryInterval() { + return recoveryInterval; + } + + public void setRecoveryInterval(int recoveryInterval) { + this.recoveryInterval = recoveryInterval; + } } diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index bda8dd3eb..359fd4d5c 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -484,6 +484,7 @@ public class KafkaMessageChannelBinder extends ? properties.getExtension().getAutoCommitOnError() : properties.getExtension().isAutoCommitOffset() && properties.getExtension().isEnableDlq(); messageListenerContainer.setAutoCommitOnError(autoCommitOnError); + messageListenerContainer.setRecoveryInterval(properties.getExtension().getRecoveryInterval()); int concurrency = Math.min(properties.getConcurrency(), listenedPartitions.size()); messageListenerContainer.setConcurrency(concurrency); diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitConsumerProperties.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitConsumerProperties.java index 2d4d5b0be..7e24b4a58 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitConsumerProperties.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitConsumerProperties.java @@ -50,6 +50,8 @@ public class RabbitConsumerProperties { private String[] replyHeaderPatterns = new String[] {"STANDARD_REPLY_HEADERS", "*"}; + private long recoveryInterval = 5000; + public String getPrefix() { return prefix; } @@ -149,4 +151,12 @@ public class RabbitConsumerProperties { public void setReplyHeaderPatterns(String[] replyHeaderPatterns) { this.replyHeaderPatterns = replyHeaderPatterns; } + + public long getRecoveryInterval() { + return recoveryInterval; + } + + public void setRecoveryInterval(long recoveryInterval) { + this.recoveryInterval = recoveryInterval; + } } diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index fd3b3ed17..a996de14d 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -280,7 +280,6 @@ public class RabbitMessageChannelBinder extends AbstractBinder 0 ? concurrency : 1; listenerContainer.setConcurrentConsumers(concurrency); @@ -288,8 +287,8 @@ public class RabbitMessageChannelBinder extends AbstractBinder concurrency) { listenerContainer.setMaxConcurrentConsumers(maxConcurrency); } - listenerContainer.setPrefetchCount(properties.getExtension().getPrefetch()); + listenerContainer.setRecoveryInterval(properties.getExtension().getRecoveryInterval()); listenerContainer.setTxSize(properties.getExtension().getTxSize()); listenerContainer.setTaskExecutor(new SimpleAsyncTaskExecutor(queue.getName() + "-")); listenerContainer.setQueues(queue); @@ -545,7 +544,7 @@ public class RabbitMessageChannelBinder extends AbstractBinder messageHeadersList) { Iterator iterator = messageHeadersList.iterator(); - Map channelsToAck = new HashMap(); + Map channelsToAck = new HashMap<>(); while (iterator.hasNext()) { MessageHeaders messageHeaders = iterator.next(); if (messageHeaders.containsKey(AmqpHeaders.CHANNEL)) { diff --git a/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc b/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc index ca2dac781..f7d32ec38 100644 --- a/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc +++ b/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc @@ -990,6 +990,10 @@ prefix:: A prefix to be added to the name of the `destination` and queues. + Default: "". +recoveryInterval:: + The interval between connection recovery attempts, in milliseconds. ++ +Default: `5000`. requeueRejected:: Whether delivery failures should be requeued. + @@ -1143,6 +1147,10 @@ If set to `true`, it will always auto-commit (if auto-commit is enabled). If not set (default), it effectively has the same value as `enableDlq`, auto-committing erroneous messages if they are sent to a DLQ, and not committing them otherwise. + Default: not set. +recoveryInterval:: + The interval between connection recovery attempts, in milliseconds. ++ +Default: `5000`. resetOffsets:: Whether to reset offsets on the consumer to the value provided by `startOffset`. +