diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpoint.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpoint.java index f1018d8363..2b3f1317ee 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpoint.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpoint.java @@ -1,5 +1,5 @@ /* - * Copyright 2013-2015 the original author or authors + * Copyright 2013-2016 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. @@ -47,6 +47,7 @@ import org.springframework.util.Assert; * @author Gunnar Hillert * @author Artem Bilan * @author Gary Russell + * @author Rainer Frey * @since 3.0 */ @ManagedResource @@ -79,6 +80,8 @@ public class RedisQueueMessageDrivenEndpoint extends MessageProducerSupport impl private volatile boolean listening; + private volatile boolean rightPop = true; + /** * @param queueName Must not be an empty String * @param connectionFactory Must not be null @@ -155,6 +158,15 @@ public class RedisQueueMessageDrivenEndpoint extends MessageProducerSupport impl this.recoveryInterval = recoveryInterval; } + /** + * Specify if {@code POP} operation from Redis List should be {@code BRPOP} or {@code BLPOP}. + * @param rightPop the {@code BRPOP} flag. Defaults to {@code true}. + * @since 4.2.5 + */ + public void setRightPop(boolean rightPop) { + this.rightPop = rightPop; + } + @Override protected void onInit() { super.onInit(); @@ -185,7 +197,12 @@ public class RedisQueueMessageDrivenEndpoint extends MessageProducerSupport impl byte[] value = null; try { - value = this.boundListOperations.rightPop(this.receiveTimeout, TimeUnit.MILLISECONDS); + if (this.rightPop) { + value = this.boundListOperations.rightPop(this.receiveTimeout, TimeUnit.MILLISECONDS); + } + else { + value = this.boundListOperations.leftPop(this.receiveTimeout, TimeUnit.MILLISECONDS); + } } catch (Exception e) { this.listening = false; @@ -224,7 +241,12 @@ public class RedisQueueMessageDrivenEndpoint extends MessageProducerSupport impl this.sendMessage(message); } else { - this.boundListOperations.rightPush(value); + if (this.rightPop) { + this.boundListOperations.rightPush(value); + } + else { + this.boundListOperations.leftPush(value); + } } } } diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/RedisQueueOutboundChannelAdapter.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/RedisQueueOutboundChannelAdapter.java index 6acff01a94..d2676a457a 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/RedisQueueOutboundChannelAdapter.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/RedisQueueOutboundChannelAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2013-2015 the original author or authors + * Copyright 2013-2016 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. @@ -33,6 +33,7 @@ import org.springframework.util.Assert; * @author Mark Fisher * @author Gunnar Hillert * @author Artem Bilan + * @author Rainer Frey * @since 3.0 */ public class RedisQueueOutboundChannelAdapter extends AbstractMessageHandler { @@ -51,6 +52,8 @@ public class RedisQueueOutboundChannelAdapter extends AbstractMessageHandler { private volatile boolean serializerExplicitlySet; + private volatile boolean leftPush = true; + public RedisQueueOutboundChannelAdapter(String queueName, RedisConnectionFactory connectionFactory) { this(new LiteralExpression(queueName), connectionFactory); } @@ -77,6 +80,15 @@ public class RedisQueueOutboundChannelAdapter extends AbstractMessageHandler { this.serializerExplicitlySet = true; } + /** + * Specify if {@code PUSH} operation to Redis List should be {@code LPUSH} or {@code RPUSH}. + * @param leftPush the {@code LPUSH} flag. Defaults to {@code true}. + * @since 4.2.5 + */ + public void setLeftPush(boolean leftPush) { + this.leftPush = leftPush; + } + public void setIntegrationEvaluationContext(EvaluationContext evaluationContext) { this.evaluationContext = evaluationContext; } @@ -113,7 +125,12 @@ public class RedisQueueOutboundChannelAdapter extends AbstractMessageHandler { } String queueName = this.queueNameExpression.getValue(this.evaluationContext, message, String.class); - this.template.boundListOps(queueName).leftPush(value); + if (this.leftPush) { + this.template.boundListOps(queueName).leftPush(value); + } + else { + this.template.boundListOps(queueName).rightPush(value); + } } }