From 31961ae5951e42bacce3f55dff2e0a0a32973045 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Fri, 29 Jan 2016 19:04:12 -0500 Subject: [PATCH] Add support for durable subscriptions in Redis --- .../stream/binder/redis/RedisMessageChannelBinder.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/RedisMessageChannelBinder.java b/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/RedisMessageChannelBinder.java index 05f705af3..1fbb70020 100644 --- a/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/RedisMessageChannelBinder.java +++ b/spring-cloud-stream-binders/spring-cloud-stream-binder-redis/src/main/java/org/springframework/cloud/stream/binder/redis/RedisMessageChannelBinder.java @@ -85,6 +85,7 @@ public class RedisMessageChannelBinder extends MessageChannelBinderSupport imple .addAll(CONSUMER_RETRY_PROPERTIES) .add(BinderPropertyKeys.CONCURRENCY) .add(BinderPropertyKeys.PARTITION_INDEX) + .add(BinderPropertyKeys.DURABLE) .build(); /** @@ -250,7 +251,10 @@ public class RedisMessageChannelBinder extends MessageChannelBinderSupport imple protected void afterUnbind(Binding binding) { if (Binding.Type.consumer.equals(binding.getType())) { String key = CONSUMER_GROUPS_KEY_PREFIX + binding.getName(); - this.redisOperations.boundZSetOps(key).incrementScore(binding.getGroup(), -1); + boolean durable = binding.getPropertiesAccessor().isDurable(defaultDurableSubscription); + if (!durable) { + this.redisOperations.boundZSetOps(key).incrementScore(binding.getGroup(), -1); + } } }