INT-3932: Redis: add rightPop and leftPush

JIRA: https://jira.spring.io/browse/INT-3932

PR review:

* remove unnecessary spaces

* fix code and doc style

* use property name rightPop/right-pop (default true)

* add reversed feature to the outbound adapter

* add new options to reference manual

* Polishing according PR comments
This commit is contained in:
Rainer Frey
2016-01-14 13:55:37 +01:00
committed by Artem Bilan
parent baa8fdc4f5
commit 84c5b7b193
13 changed files with 206 additions and 23 deletions

View File

@@ -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();

View File

@@ -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();
}

View File

@@ -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);
}
}
}
}

View File

@@ -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);
}
}
}

View File

@@ -390,6 +390,14 @@
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="right-pop" type="xsd:string" default="true">
<xsd:annotation>
<xsd:documentation>
When 'false', specifies that data is read using a 'left pop' operation instead of a 'right pop'.
Default is 'true'.
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="task-executor" type="xsd:string">
<xsd:annotation>
<xsd:documentation><![CDATA[
@@ -605,7 +613,17 @@
<xsd:attribute name="extract-payload" type="xsd:string" default="true">
<xsd:annotation>
<xsd:documentation>
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'.
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="left-push" type="xsd:string" default="true">
<xsd:annotation>
<xsd:documentation>
When 'false', specifies that data is written using a 'right push' operation instead
of a 'left push'.
Default is 'true'.
</xsd:documentation>
</xsd:annotation>

View File

@@ -31,7 +31,8 @@
recovery-interval="3000"
task-executor="executor"
auto-startup="false"
phase="100"/>
phase="100"
right-pop="false"/>
<bean id="executor" class="org.springframework.integration.util.ErrorHandlingTaskExecutor">
<constructor-arg ref="threadPoolTaskExecutor"/>

View File

@@ -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));
}
}

View File

@@ -26,7 +26,8 @@
queue-expression="headers['redis_queue']"
extract-payload="false"
serializer="serializer"
connection-factory="customRedisConnectionFactory"/>
connection-factory="customRedisConnectionFactory"
left-push="false"/>
<bean id="serializer" class="org.springframework.data.redis.serializer.StringRedisSerializer"/>

View File

@@ -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));
}
}

View File

@@ -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<String, Object> redisTemplate = new RedisTemplate<String, Object>();
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<Object> receive = (Message<Object>) channel.receive(2000);
assertNotNull(receive);
assertEquals(payload, receive.getPayload());
receive = (Message<Object>) channel.receive(2000);
assertNotNull(receive);
assertEquals(payload2, receive.getPayload());
endpoint.stop();
}
private void waitListening(RedisQueueMessageDrivenEndpoint endpoint) throws InterruptedException {
int n = 0;
do {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2013-2014 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.
@@ -49,6 +49,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
/**
* @author Gunnar Hillert
* @author Artem Bilan
* @author Rainer Frey
* @since 3.0
*/
@ContextConfiguration
@@ -177,4 +178,42 @@ public class RedisQueueOutboundChannelAdapterTests extends RedisAvailableTests {
assertEquals(message.getPayload(), resultMessage.getPayload());
}
@Test
@RedisAvailable
public void testInt3932LeftPushFalse() throws Exception {
final String queueName = "si.test.Int3932LeftPushFalse";
final RedisQueueOutboundChannelAdapter handler = new RedisQueueOutboundChannelAdapter(queueName,
this.connectionFactory);
handler.setLeftPush(false);
String payload = "testing";
handler.handleMessage(MessageBuilder.withPayload(payload).build());
Date payload2 = new Date();
handler.handleMessage(MessageBuilder.withPayload(payload2).build());
RedisTemplate<String, ?> redisTemplate = new StringRedisTemplate();
redisTemplate.setConnectionFactory(this.connectionFactory);
redisTemplate.afterPropertiesSet();
Object result = redisTemplate.boundListOps(queueName).leftPop(5000, TimeUnit.MILLISECONDS);
assertNotNull(result);
assertEquals(payload, result);
RedisTemplate<String, ?> redisTemplate2 = new RedisTemplate<String, Object>();
redisTemplate2.setConnectionFactory(this.connectionFactory);
redisTemplate2.setEnableDefaultSerializer(false);
redisTemplate2.setKeySerializer(new StringRedisSerializer());
redisTemplate2.setValueSerializer(new JdkSerializationRedisSerializer());
redisTemplate2.afterPropertiesSet();
Object result2 = redisTemplate2.boundListOps(queueName).leftPop(5000, TimeUnit.MILLISECONDS);
assertNotNull(result2);
assertEquals(payload2, result2);
}
}