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