From 42cc18e4619eb70ebd6db69ca93c20c266599985 Mon Sep 17 00:00:00 2001 From: David Turanski Date: Thu, 25 Aug 2011 17:40:33 -0400 Subject: [PATCH] INT-2082 Added supported event types --- .../ContinuousQueryMessageProducer.java | 54 ++++++++++++++----- .../gemfire/inbound/CqEventType.java | 30 +++++++++++ ...nuousQueryMessageProducerTests-context.xml | 2 +- 3 files changed, 73 insertions(+), 13 deletions(-) create mode 100644 spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CqEventType.java 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 db0a823037..b5727ea112 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,6 +16,10 @@ package org.springframework.integration.gemfire.inbound; +import java.util.Arrays; +import java.util.HashSet; +import java.util.Set; + import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.data.gemfire.listener.CqQueryDefinition; @@ -40,15 +44,22 @@ import com.gemstone.gemfire.cache.query.CqEvent; */ public class ContinuousQueryMessageProducer extends SpelMessageProducerSupport implements QueryListener { private static Log logger = LogFactory.getLog(ContinuousQueryMessageProducer.class); - + private final String query; + private final QueryListenerContainer queryListenerContainer; + private volatile String queryName; + private boolean durable; - + + private volatile Set supportedEventTypes = new HashSet(Arrays.asList(CqEventType.CREATED, + CqEventType.UPDATED)); + /** * - * @param queryListenerContainer a {@link org.springframework.data.gemfire.listener.QueryListenerContainer} + * @param queryListenerContainer a + * {@link org.springframework.data.gemfire.listener.QueryListenerContainer} * @param query the query string */ public ContinuousQueryMessageProducer(QueryListenerContainer queryListenerContainer, String query) { @@ -57,7 +68,7 @@ public class ContinuousQueryMessageProducer extends SpelMessageProducerSupport i this.queryListenerContainer = queryListenerContainer; this.query = query; } - + /** * * @param queryName optional query name @@ -65,7 +76,7 @@ public class ContinuousQueryMessageProducer extends SpelMessageProducerSupport i public void setQueryName(String queryName) { this.queryName = queryName; } - + /** * * @param durable true if the query is a durable subscription @@ -74,16 +85,22 @@ public class ContinuousQueryMessageProducer extends SpelMessageProducerSupport i this.durable = durable; } + public void setSupportedEventTypes(CqEventType... eventTypes) { + Assert.notEmpty(eventTypes, "eventTypes must not be empty"); + this.supportedEventTypes = new HashSet(Arrays.asList(eventTypes)); + } + @Override protected void onInit() { super.onInit(); - if (queryName == null){ + if (queryName == null) { queryListenerContainer.addListener(new CqQueryDefinition(this.query, this, this.durable)); - } else { + } + else { queryListenerContainer.addListener(new CqQueryDefinition(this.queryName, this.query, this, this.durable)); } } - + /* * (non-Javadoc) * @@ -92,11 +109,24 @@ public class ContinuousQueryMessageProducer extends SpelMessageProducerSupport i * .gemfire.cache.query.CqEvent) */ public void onEvent(CqEvent event) { - if (logger.isDebugEnabled()){ - logger.debug(String.format("processing cq event key [%s] event [%s]",event.getBaseOperation().toString(),event.getKey())); + if (isEventSupported(event)) { + if (logger.isDebugEnabled()) { + logger.debug(String.format("processing cq event key [%s] event [%s]", event.getBaseOperation() + .toString(), event.getKey())); + } + Message cqEventMessage = MessageBuilder.withPayload(evaluationResult(event)).build(); + sendMessage(cqEventMessage); } - Message cqEventMessage = MessageBuilder.withPayload(evaluationResult(event)).build(); - sendMessage(cqEventMessage); + } + + /** + * @param event + * @return + */ + private boolean isEventSupported(CqEvent event) { + String eventName = event.getBaseOperation().toString()+"D"; + CqEventType eventType = CqEventType.valueOf(eventName); + return supportedEventTypes.contains(eventType); } } \ No newline at end of file diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CqEventType.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CqEventType.java new file mode 100644 index 0000000000..1401879371 --- /dev/null +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CqEventType.java @@ -0,0 +1,30 @@ +/* + * 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; + +/** + * Enumeration of GemFire Continuous Query Event Types + * @author David Turanski + * @since 2.1 + */ +public enum CqEventType { + CREATED, + + UPDATED, + + DESTROYED, + + REGION_CLEARED, + + REGION_INVALIDATED +} 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 761d86c414..60665d582b 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 @@ -5,7 +5,7 @@ xmlns:int="http://www.springframework.org/schema/integration" xmlns:util="http://www.springframework.org/schema/util" xmlns:int-gfe="http://www.springframework.org/schema/integration/gemfire" - xsi:schemaLocation="http://www.springframework.org/schema/integration/gemfire http://www.springframework.org/schema/integration/scripting/spring-integration-gemfire.xsd + xsi:schemaLocation="http://www.springframework.org/schema/integration/gemfire http://www.springframework.org/schema/integration/gemfire/spring-integration-gemfire.xsd http://www.springframework.org/schema/gemfire http://www.springframework.org/schema/gemfire/spring-gemfire-1.1.xsd http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd