From 2cade4220e402e4c5f2833970958727fd2b028ef Mon Sep 17 00:00:00 2001 From: David Turanski Date: Thu, 25 Aug 2011 13:06:00 -0400 Subject: [PATCH] INT-2082 added support for evaluating payload expression in ContinuousQueryMessageProducer --- .../CacheListeningMessageProducer.java | 31 ++--------- .../ContinuousQueryMessageProducer.java | 10 ++-- .../inbound/SpelMessageProducerSupport.java | 54 +++++++++++++++++++ .../GemfireInboundChannelAdapterTests.java | 1 - ...nuousQueryMessageProducerTests-context.xml | 18 +++++-- .../ContinuousQueryMessageProducerTests.java | 45 ++++++++++------ 6 files changed, 108 insertions(+), 51 deletions(-) create mode 100644 spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/SpelMessageProducerSupport.java diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducer.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducer.java index 8d7296e8c0..3cbf6c97e2 100644 --- a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducer.java +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducer.java @@ -22,10 +22,6 @@ import java.util.Set; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; - -import org.springframework.expression.Expression; -import org.springframework.expression.spel.standard.SpelExpressionParser; -import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.integration.support.MessageBuilder; import org.springframework.util.Assert; @@ -47,7 +43,7 @@ import com.gemstone.gemfire.cache.util.CacheListenerAdapter; * @since 2.1 */ @SuppressWarnings({"rawtypes", "unchecked"}) -public class CacheListeningMessageProducer extends MessageProducerSupport { +public class CacheListeningMessageProducer extends SpelMessageProducerSupport { private final Log logger = LogFactory.getLog(this.getClass()); @@ -58,9 +54,6 @@ public class CacheListeningMessageProducer extends MessageProducerSupport { private volatile Set supportedEventTypes = new HashSet(Arrays.asList(EventType.CREATED, EventType.UPDATED)); - private volatile Expression payloadExpression; - - private final SpelExpressionParser parser = new SpelExpressionParser(); public CacheListeningMessageProducer(Region region) { Assert.notNull(region, "region must not be null"); @@ -74,14 +67,6 @@ public class CacheListeningMessageProducer extends MessageProducerSupport { this.supportedEventTypes = new HashSet(Arrays.asList(eventTypes)); } - public void setPayloadExpression(String payloadExpression) { - if (payloadExpression == null) { - this.payloadExpression = null; - } - else { - this.payloadExpression = this.parser.parseExpression(payloadExpression); - } - } @Override protected void doStart() { @@ -105,8 +90,7 @@ public class CacheListeningMessageProducer extends MessageProducerSupport { } } - - + private class MessageProducingCacheListener extends CacheListenerAdapter { @Override @@ -137,14 +121,9 @@ public class CacheListeningMessageProducer extends MessageProducerSupport { } } - private void processEvent(EntryEvent event) { - if (payloadExpression != null) { - Object evaluationResult = payloadExpression.getValue(event); - this.publish(evaluationResult); - } - else { - this.publish(event); - } + private void processEvent(EntryEvent event) { + this.publish(evaluationResult(event)); + } private void publish(Object payload) { diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducer.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducer.java index a16b38f26b..db0a823037 100644 --- a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducer.java +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducer.java @@ -16,15 +16,15 @@ package org.springframework.integration.gemfire.inbound; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.springframework.data.gemfire.listener.CqQueryDefinition; import org.springframework.data.gemfire.listener.QueryListener; import org.springframework.data.gemfire.listener.QueryListenerContainer; import org.springframework.integration.Message; -import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.integration.support.MessageBuilder; import org.springframework.util.Assert; -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; + import com.gemstone.gemfire.cache.query.CqEvent; /** @@ -38,7 +38,7 @@ import com.gemstone.gemfire.cache.query.CqEvent; * @since 2.1 * */ -public class ContinuousQueryMessageProducer extends MessageProducerSupport implements QueryListener { +public class ContinuousQueryMessageProducer extends SpelMessageProducerSupport implements QueryListener { private static Log logger = LogFactory.getLog(ContinuousQueryMessageProducer.class); private final String query; @@ -95,7 +95,7 @@ public class ContinuousQueryMessageProducer extends MessageProducerSupport imple if (logger.isDebugEnabled()){ logger.debug(String.format("processing cq event key [%s] event [%s]",event.getBaseOperation().toString(),event.getKey())); } - Message cqEventMessage = MessageBuilder.withPayload(event).build(); + Message cqEventMessage = MessageBuilder.withPayload(evaluationResult(event)).build(); sendMessage(cqEventMessage); } diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/SpelMessageProducerSupport.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/SpelMessageProducerSupport.java new file mode 100644 index 0000000000..c4332dea9d --- /dev/null +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/SpelMessageProducerSupport.java @@ -0,0 +1,54 @@ +/* + * Copyright 2002-2011 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. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on + * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the + * specific language governing permissions and limitations under the License. + */ +package org.springframework.integration.gemfire.inbound; + +import org.springframework.expression.Expression; +import org.springframework.expression.spel.standard.SpelExpressionParser; +import org.springframework.integration.endpoint.MessageProducerSupport; + +/** + * @author David Turanski + * @since 2.1 + * + */ +public class SpelMessageProducerSupport extends MessageProducerSupport { + + private volatile Expression payloadExpression; + + private final SpelExpressionParser parser = new SpelExpressionParser(); + + + @Override + protected void onInit(){ + super.onInit(); + } + + public void setPayloadExpression(String payloadExpression) { + if (payloadExpression == null) { + this.payloadExpression = null; + } + else { + this.payloadExpression = this.parser.parseExpression(payloadExpression); + } + } + + protected Object evaluationResult(Object payload){ + Object evaluationResult = payload; + if (payloadExpression != null) { + evaluationResult = payloadExpression.getValue(payload); + } + return evaluationResult; + } + + +} diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/GemfireInboundChannelAdapterTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/GemfireInboundChannelAdapterTests.java index 5201d86214..003bac50be 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/GemfireInboundChannelAdapterTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/GemfireInboundChannelAdapterTests.java @@ -23,7 +23,6 @@ import org.springframework.integration.MessagingException; import org.springframework.integration.core.MessageHandler; import org.springframework.integration.core.SubscribableChannel; import org.springframework.integration.message.ErrorMessage; -import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/cq/ContinuousQueryMessageProducerTests-context.xml b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/cq/ContinuousQueryMessageProducerTests-context.xml index 0b5c54e67a..b5cca21c19 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/cq/ContinuousQueryMessageProducerTests-context.xml +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/cq/ContinuousQueryMessageProducerTests-context.xml @@ -23,14 +23,26 @@ - + - + - + + + + + + + + + + + + + diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/cq/ContinuousQueryMessageProducerTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/cq/ContinuousQueryMessageProducerTests.java index 7e6640fa1e..c87de95792 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/cq/ContinuousQueryMessageProducerTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/cq/ContinuousQueryMessageProducerTests.java @@ -12,6 +12,7 @@ */ package org.springframework.integration.gemfire.inbound.cq; +import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertTrue; @@ -19,6 +20,7 @@ import java.io.IOException; import java.io.OutputStream; import org.junit.AfterClass; +import org.junit.Before; import org.junit.BeforeClass; import org.junit.Test; import org.junit.runner.RunWith; @@ -41,6 +43,8 @@ import com.gemstone.gemfire.internal.cache.LocalRegion; @ContextConfiguration public class ContinuousQueryMessageProducerTests { + static ConfigurableApplicationContext staticCtx; + @Autowired LocalRegion region; @@ -48,7 +52,10 @@ public class ContinuousQueryMessageProducerTests { ConfigurableApplicationContext applicationContext; @Autowired - PollableChannel outputChannel; + PollableChannel outputChannel1; + + @Autowired + PollableChannel outputChannel2; static OutputStream os; @BeforeClass @@ -56,30 +63,36 @@ public class ContinuousQueryMessageProducerTests { os = ForkUtil.cacheServer(); } - + + @Before + public void setUp() { + staticCtx = applicationContext; + } @Test - public void test() throws InterruptedException { + public void testCqEvent() throws InterruptedException { region.put("one",1); - Message msg = outputChannel.receive(1000); + Message msg = outputChannel1.receive(1000); assertNotNull(msg); assertTrue(msg.getPayload() instanceof CqEvent); - /* - * Avoid shutdown errors - */ - applicationContext.close(); + } + + @Test + public void testPayloadExpression() throws InterruptedException { + region.put("one",1); + Message msg = outputChannel2.receive(1000); + assertNotNull(msg); + assertEquals(1,msg.getPayload()); + + } @AfterClass public static void cleanUp() { - - try { - Thread.sleep(3000); - } - catch (InterruptedException e) { - // TODO Auto-generated catch block - e.printStackTrace(); - } + /* + * Avoid shutdown errors + */ + staticCtx.close(); sendSignal(); }