diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParser.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParser.java
index 401b413dea..7e1daf2026 100644
--- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParser.java
+++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParser.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013 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.
@@ -30,6 +30,7 @@ import org.springframework.util.StringUtils;
* Parser for the <queue-inbound-channel-adapter> element of the 'redis' namespace.
*
* @author Artem Bilan
+ * @author Rainer Frey
* @since 3.0
*/
public class RedisQueueInboundChannelAdapterParser extends AbstractChannelAdapterParser {
@@ -51,6 +52,7 @@ public class RedisQueueInboundChannelAdapterParser extends AbstractChannelAdapte
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "expect-message");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "receive-timeout");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "recovery-interval");
+ IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "right-pop");
builder.addPropertyReference("outputChannel", channelName);
return builder.getBeanDefinition();
diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParser.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParser.java
index 9e0653e814..6eb85a8242 100644
--- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParser.java
+++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParser.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2013 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.
@@ -31,6 +31,7 @@ import org.springframework.util.StringUtils;
* Parser for the <int-redis:queue-outbound-channel-adapter> element.
*
* @author Artem Bilan
+ * @author Rainer Frey
* @since 3.0
*/
public class RedisQueueOutboundChannelAdapterParser extends AbstractOutboundChannelAdapterParser {
@@ -50,6 +51,7 @@ public class RedisQueueOutboundChannelAdapterParser extends AbstractOutboundChan
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "extract-payload");
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "serializer");
+ IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "left-push");
return builder.getBeanDefinition();
}
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 fb72d397f2..6e9eeaea68 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
@@ -48,6 +48,7 @@ import org.springframework.util.Assert;
* @author Gunnar Hillert
* @author Artem Bilan
* @author Gary Russell
+ * @author Rainer Frey
* @since 3.0
*/
@ManagedResource
@@ -80,6 +81,8 @@ public class RedisQueueMessageDrivenEndpoint extends MessageProducerSupport impl
private volatile Runnable stopCallback;
+ private volatile boolean rightPop = true;
+
/**
* @param queueName Must not be an empty String
* @param connectionFactory Must not be null
@@ -158,6 +161,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.3
+ */
+ public void setRightPop(boolean rightPop) {
+ this.rightPop = rightPop;
+ }
+
@Override
protected void onInit() {
super.onInit();
@@ -188,7 +200,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;
@@ -227,7 +244,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..14d47d540a 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.3
+ */
+ 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);
+ }
}
}
diff --git a/spring-integration-redis/src/main/resources/org/springframework/integration/redis/config/spring-integration-redis-4.3.xsd b/spring-integration-redis/src/main/resources/org/springframework/integration/redis/config/spring-integration-redis-4.3.xsd
index 90405da5d5..ec78f2268e 100644
--- a/spring-integration-redis/src/main/resources/org/springframework/integration/redis/config/spring-integration-redis-4.3.xsd
+++ b/spring-integration-redis/src/main/resources/org/springframework/integration/redis/config/spring-integration-redis-4.3.xsd
@@ -390,6 +390,14 @@
+
+
+
+ When 'false', specifies that data is read using a 'left pop' operation instead of a 'right pop'.
+ Default is 'true'.
+
+
+
- Specifies if the Message payload or the entire (serialized) Message will be send to the Redis queue.
+ Specifies if the Message payload or the entire (serialized) Message will be send
+ to the Redis queue.
+ Default is 'true'.
+
+
+
+
+
+
+ When 'false', specifies that data is written using a 'right push' operation instead
+ of a 'left push'.
Default is 'true'.
diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParserTests-context.xml b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParserTests-context.xml
index 3a1cc5b604..b149ac3e6b 100644
--- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParserTests-context.xml
+++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParserTests-context.xml
@@ -31,7 +31,8 @@
recovery-interval="3000"
task-executor="executor"
auto-startup="false"
- phase="100"/>
+ phase="100"
+ right-pop="false"/>
diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParserTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParserTests.java
index 7376a7fd45..213a4eb587 100644
--- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParserTests.java
+++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueInboundChannelAdapterParserTests.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.
@@ -44,6 +44,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
/**
* @author Artem Bilan
* @author Gary Russell
+ * @author Rainer Frey
* @since 3.0
*/
@ContextConfiguration
@@ -103,6 +104,7 @@ public class RedisQueueInboundChannelAdapterParserTests {
assertTrue(TestUtils.getPropertyValue(this.defaultAdapter, "autoStartup", Boolean.class));
assertEquals(Integer.MAX_VALUE / 2, TestUtils.getPropertyValue(this.defaultAdapter, "phase"));
assertSame(this.defaultAdapterChannel, TestUtils.getPropertyValue(this.defaultAdapter, "outputChannel"));
+ assertTrue(TestUtils.getPropertyValue(this.defaultAdapter, "rightPop", Boolean.class));
}
@Test
@@ -119,6 +121,7 @@ public class RedisQueueInboundChannelAdapterParserTests {
assertFalse(TestUtils.getPropertyValue(this.customAdapter, "autoStartup", Boolean.class));
assertEquals(100, TestUtils.getPropertyValue(this.customAdapter, "phase"));
assertSame(this.sendChannel, TestUtils.getPropertyValue(this.customAdapter, "outputChannel"));
+ assertFalse(TestUtils.getPropertyValue(this.customAdapter, "rightPop", Boolean.class));
}
}
diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParserTests-context.xml b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParserTests-context.xml
index 0e585480c1..749ca20573 100644
--- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParserTests-context.xml
+++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParserTests-context.xml
@@ -26,7 +26,8 @@
queue-expression="headers['redis_queue']"
extract-payload="false"
serializer="serializer"
- connection-factory="customRedisConnectionFactory"/>
+ connection-factory="customRedisConnectionFactory"
+ left-push="false"/>
diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParserTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParserTests.java
index e111d4489f..47d7f039c8 100644
--- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParserTests.java
+++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/config/RedisQueueOutboundChannelAdapterParserTests.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.
@@ -44,6 +44,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
/**
* @author Artem Bilan
* @author Gary Russell
+ * @author Rainer Frey
* @since 3.0
*/
@ContextConfiguration
@@ -89,6 +90,7 @@ public class RedisQueueOutboundChannelAdapterParserTests {
assertThat(TestUtils.getPropertyValue(handler, "h.advised.advisors.first.item.advice"),
Matchers.instanceOf(RequestHandlerRetryAdvice.class));
+ assertTrue(TestUtils.getPropertyValue(this.defaultAdapter, "leftPush", Boolean.class));
}
@Test
@@ -98,6 +100,7 @@ public class RedisQueueOutboundChannelAdapterParserTests {
assertFalse(TestUtils.getPropertyValue(this.customAdapter, "extractPayload", Boolean.class));
assertTrue(TestUtils.getPropertyValue(this.customAdapter, "serializerExplicitlySet", Boolean.class));
assertSame(this.serializer, TestUtils.getPropertyValue(this.customAdapter, "serializer"));
+ assertFalse(TestUtils.getPropertyValue(this.customAdapter, "leftPush", Boolean.class));
}
}
diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java
index aae88e6ccb..0d9bb8caad 100644
--- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java
+++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/RedisQueueMessageDrivenEndpointTests.java
@@ -73,6 +73,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
* @author Gunnar Hillert
* @author Artem Bilan
* @author Gary Russell
+ * @author Rainer Frey
* @since 3.0
*/
@ContextConfiguration
@@ -157,7 +158,8 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
PollableChannel errorChannel = new QueueChannel();
- RedisQueueMessageDrivenEndpoint endpoint = new RedisQueueMessageDrivenEndpoint(queueName, this.connectionFactory);
+ RedisQueueMessageDrivenEndpoint endpoint =
+ new RedisQueueMessageDrivenEndpoint(queueName, this.connectionFactory);
endpoint.setBeanFactory(Mockito.mock(BeanFactory.class));
endpoint.setExpectMessage(true);
endpoint.setOutputChannel(channel);
@@ -175,7 +177,8 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
assertNotNull(receive);
assertThat(receive, Matchers.instanceOf(ErrorMessage.class));
assertThat(receive.getPayload(), Matchers.instanceOf(MessagingException.class));
- assertThat(((Exception) receive.getPayload()).getMessage(), Matchers.containsString("Deserialization of Message failed."));
+ assertThat(((Exception) receive.getPayload()).getMessage(),
+ Matchers.containsString("Deserialization of Message failed."));
assertThat(((Exception) receive.getPayload()).getCause(), Matchers.instanceOf(ClassCastException.class));
assertThat(((Exception) receive.getPayload()).getCause().getMessage(),
Matchers.containsString("java.lang.String cannot be cast to org.springframework.messaging.Message"));
@@ -193,7 +196,8 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
redisTemplate.setConnectionFactory(this.connectionFactory);
redisTemplate.afterPropertiesSet();
- redisTemplate.boundListOps("si.test.Int3017IntegrationInbound").leftPush("{\"payload\":\"" + payload + "\",\"headers\":{}}");
+ redisTemplate.boundListOps("si.test.Int3017IntegrationInbound")
+ .leftPush("{\"payload\":\"" + payload + "\",\"headers\":{}}");
Message> receive = this.fromChannel.receive(2000);
assertNotNull(receive);
@@ -340,6 +344,50 @@ public class RedisQueueMessageDrivenEndpointTests extends RedisAvailableTests {
endpoint.stop();
}
+ @Test
+ @RedisAvailable
+ @SuppressWarnings("unchecked")
+ public void testInt3932ReadFromLeft() throws Exception {
+
+ String queueName = "si.test.redisQueueInboundChannelAdapterTests3932";
+
+ RedisTemplate redisTemplate = new RedisTemplate();
+ redisTemplate.setConnectionFactory(this.connectionFactory);
+ redisTemplate.setEnableDefaultSerializer(false);
+ redisTemplate.setKeySerializer(new StringRedisSerializer());
+ redisTemplate.setValueSerializer(new JdkSerializationRedisSerializer());
+ redisTemplate.afterPropertiesSet();
+
+ String payload = "testing";
+
+ redisTemplate.boundListOps(queueName).rightPush(payload);
+
+ Date payload2 = new Date();
+
+ redisTemplate.boundListOps(queueName).rightPush(payload2);
+
+ PollableChannel channel = new QueueChannel();
+
+ RedisQueueMessageDrivenEndpoint endpoint =
+ new RedisQueueMessageDrivenEndpoint(queueName, this.connectionFactory);
+ endpoint.setBeanFactory(Mockito.mock(BeanFactory.class));
+ endpoint.setOutputChannel(channel);
+ endpoint.setReceiveTimeout(1000);
+ endpoint.setRightPop(false);
+ endpoint.afterPropertiesSet();
+ endpoint.start();
+
+ Message