From fd9b74b41cfb0a3c13194c874e1a994cae242ad3 Mon Sep 17 00:00:00 2001 From: Costin Leau Date: Tue, 25 Jan 2011 13:33:24 +0200 Subject: [PATCH] DATAKV-25 + add error handler setter on the container + improve docs --- .../config/RedisListenerContainerParser.java | 11 ++ .../RedisMessageListenerContainer.java | 115 +++++++++++++++--- .../redis/config/spring-redis-1.0.xsd | 13 -- .../keyvalue/redis/config/NamespaceTest.java | 13 ++ .../redis/config/StubErrorHandler.java | 35 ++++++ .../adapter/ThrowableMessageListener.java | 31 +++++ .../data/keyvalue/redis/config/namespace.xml | 6 +- src/docbkx/appendix/appendix-schema.xml | 16 +++ src/docbkx/appendix/introduction.xml | 10 ++ src/docbkx/index.xml | 7 ++ src/docbkx/reference/redis-messaging.xml | 23 +++- src/docbkx/reference/redis.xml | 6 +- 12 files changed, 252 insertions(+), 34 deletions(-) create mode 100644 spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/config/StubErrorHandler.java create mode 100644 spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/listener/adapter/ThrowableMessageListener.java create mode 100644 src/docbkx/appendix/appendix-schema.xml create mode 100644 src/docbkx/appendix/introduction.xml diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/config/RedisListenerContainerParser.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/config/RedisListenerContainerParser.java index 8179a4237..1c300a8f4 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/config/RedisListenerContainerParser.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/config/RedisListenerContainerParser.java @@ -63,6 +63,12 @@ class RedisListenerContainerParser extends AbstractSimpleBeanDefinitionParser { builder.addPropertyReference(propertyName, attribute.getValue()); } } + + String phase = element.getAttribute("phase"); + if (StringUtils.hasText(phase)) { + builder.addPropertyValue("phase", phase); + } + postProcess(builder, element); // parse nested listeners @@ -81,6 +87,11 @@ class RedisListenerContainerParser extends AbstractSimpleBeanDefinitionParser { } } + @Override + protected boolean isEligibleAttribute(String attributeName) { + return (!"phase".equals(attributeName)); + } + /** * Parses a listener definition. Returns the listener bean reference definition (as the array first entry) and its associated topics (also as bean definitions). * diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/listener/RedisMessageListenerContainer.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/listener/RedisMessageListenerContainer.java index f2620444d..5cf29e489 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/listener/RedisMessageListenerContainer.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/listener/RedisMessageListenerContainer.java @@ -44,6 +44,7 @@ import org.springframework.data.keyvalue.redis.serializer.StringRedisSerializer; import org.springframework.scheduling.SchedulingAwareRunnable; import org.springframework.util.ClassUtils; import org.springframework.util.CollectionUtils; +import org.springframework.util.ErrorHandler; /** * Container providing asynchronous behaviour for Redis message listeners. @@ -60,7 +61,10 @@ import org.springframework.util.CollectionUtils; */ public class RedisMessageListenerContainer implements InitializingBean, DisposableBean, BeanNameAware, SmartLifecycle { - private static final Log log = LogFactory.getLog(RedisMessageListenerContainer.class); + /** Logger available to subclasses */ + protected final Log logger = LogFactory.getLog(getClass()); + + /** * Default thread name prefix: "RedisListeningContainer-". @@ -78,6 +82,8 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab private String beanName; + private ErrorHandler errorHandler; + private final Object monitor = new Object(); // whether the container is running (or not) @@ -115,8 +121,9 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab subscriptionExecutor = taskExecutor; } - start(); initialized = true; + + start(); } /** @@ -140,8 +147,8 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab if (taskExecutor instanceof DisposableBean) { ((DisposableBean) taskExecutor).destroy(); - if (log.isDebugEnabled()) { - log.debug("Stopped internally-managed task executor"); + if (logger.isDebugEnabled()) { + logger.debug("Stopped internally-managed task executor"); } } } @@ -185,8 +192,8 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab } } - if (log.isDebugEnabled()) { - log.debug("Started RedisMessageListenerContainer"); + if (logger.isDebugEnabled()) { + logger.debug("Started RedisMessageListenerContainer"); } } } @@ -207,8 +214,74 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab } } - if (log.isDebugEnabled()) { - log.debug("Stopped RedisMessageListenerContainer"); + if (logger.isDebugEnabled()) { + logger.debug("Stopped RedisMessageListenerContainer"); + } + } + + + /** + * Process a message received from the provider. + * + * @param message + * @param pattern + */ + protected void processMessage(MessageListener listener, Message message, byte[] pattern) { + executeListener(listener, message, pattern); + } + + + /** + * Execute the specified listener. + * + * @see #handleListenerException + */ + protected void executeListener(MessageListener listener, Message message, byte[] pattern) { + try { + listener.onMessage(message, pattern); + } catch (Throwable ex) { + handleListenerException(ex); + } + } + + /** + * Return whether this container is currently active, + * that is, whether it has been set up but not shut down yet. + */ + public final boolean isActive() { + return initialized; + } + + /** + * Handle the given exception that arose during listener execution. + *

The default implementation logs the exception at error level. + * This can be overridden in subclasses. + * @param ex the exception to handle + */ + protected void handleListenerException(Throwable ex) { + if (isActive()) { + // Regular case: failed while active. + // Invoke ErrorHandler if available. + invokeErrorHandler(ex); + } + else { + // Rare case: listener thread failed after container shutdown. + // Log at debug level, to avoid spamming the shutdown logger. + logger.debug("Listener exception after container shutdown", ex); + } + } + + /** + * Invoke the registered ErrorHandler, if any. Log at error level otherwise. + * @param ex the uncaught error that arose during message processing. + * @see #setErrorHandler + */ + protected void invokeErrorHandler(Throwable ex) { + if (this.errorHandler != null) { + this.errorHandler.handleError(ex); + } + else if (logger.isWarnEnabled()) { + logger.warn("Execution of JMS message listener failed, and no ErrorHandler has been set.", ex); } } @@ -270,6 +343,15 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab this.serializer = serializer; } + /** + * Set an ErrorHandler to be invoked in case of any uncaught exceptions thrown + * while processing a Message. By default there will be no ErrorHandler + * so that error-level logging is the only result. + */ + public void setErrorHandler(ErrorHandler errorHandler) { + this.errorHandler = errorHandler; + } + /** * Attaches the given listeners (and their topics) to the container. * @@ -332,7 +414,7 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab * Method inspecting whether listening for messages (and thus using a thread) is actually needed and triggering it. */ private void lazyListen() { - boolean debug = log.isDebugEnabled(); + boolean debug = logger.isDebugEnabled(); boolean started = false; if (isRunning()) { @@ -348,10 +430,10 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab } if (debug) { if (started) { - log.debug("Started listening for Redis messages"); + logger.debug("Started listening for Redis messages"); } else { - log.debug("Postpone listening for Redis messages until actual listeners are added"); + logger.debug("Postpone listening for Redis messages until actual listeners are added"); } } } @@ -362,7 +444,7 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab List channels = new ArrayList(topics.size()); List patterns = new ArrayList(topics.size()); - boolean trace = log.isTraceEnabled(); + boolean trace = logger.isTraceEnabled(); for (Topic topic : topics) { @@ -378,7 +460,7 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab channels.add(holder.array); if (trace) - log.trace("Adding listener '" + listener + "' on channel '" + topic.getTopic() + "'"); + logger.trace("Adding listener '" + listener + "' on channel '" + topic.getTopic() + "'"); } else if (topic instanceof PatternTopic) { @@ -391,7 +473,7 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab patterns.add(holder.array); if (trace) - log.trace("Adding listener '" + listener + "' for pattern '" + topic.getTopic() + "'"); + logger.trace("Adding listener '" + listener + "' for pattern '" + topic.getTopic() + "'"); } else { @@ -406,6 +488,7 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab } } + /** * Runnable used for Redis subscription. Implemented as a dedicated class to provide as many hints * as possible to the underlying thread pool. @@ -639,7 +722,7 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab taskExecutor.execute(new Runnable() { @Override public void run() { - messageListener.onMessage(message, null); + processMessage(messageListener, message, null); } }); } @@ -650,7 +733,7 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab taskExecutor.execute(new Runnable() { @Override public void run() { - messageListener.onMessage(message, pattern.clone()); + processMessage(messageListener, message, pattern.clone()); } }); } diff --git a/spring-data-redis/src/main/resources/org/springframework/data/keyvalue/redis/config/spring-redis-1.0.xsd b/spring-data-redis/src/main/resources/org/springframework/data/keyvalue/redis/config/spring-redis-1.0.xsd index e0cd544a3..0ba5d1ea9 100644 --- a/spring-data-redis/src/main/resources/org/springframework/data/keyvalue/redis/config/spring-redis-1.0.xsd +++ b/spring-data-redis/src/main/resources/org/springframework/data/keyvalue/redis/config/spring-redis-1.0.xsd @@ -84,19 +84,6 @@ - - - - - - - - - - throwables = new LinkedBlockingDeque(); + + @Override + public void handleError(Throwable t) { + throwables.add(t); + } + +} diff --git a/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/listener/adapter/ThrowableMessageListener.java b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/listener/adapter/ThrowableMessageListener.java new file mode 100644 index 000000000..2f2903e22 --- /dev/null +++ b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/listener/adapter/ThrowableMessageListener.java @@ -0,0 +1,31 @@ +/* + * Copyright 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.data.keyvalue.redis.listener.adapter; + +import org.springframework.data.keyvalue.redis.connection.Message; +import org.springframework.data.keyvalue.redis.connection.MessageListener; + +/** + * + * @author Costin Leau + */ +public class ThrowableMessageListener implements MessageListener { + + @Override + public void onMessage(Message message, byte[] pattern) { + throw new IllegalStateException("throwing exception for message " + message); + } +} diff --git a/spring-data-redis/src/test/resources/org/springframework/data/keyvalue/redis/config/namespace.xml b/spring-data-redis/src/test/resources/org/springframework/data/keyvalue/redis/config/namespace.xml index 2ae0ace02..91a7cc2f3 100644 --- a/spring-data-redis/src/test/resources/org/springframework/data/keyvalue/redis/config/namespace.xml +++ b/spring-data-redis/src/test/resources/org/springframework/data/keyvalue/redis/config/namespace.xml @@ -12,17 +12,21 @@ - + + + + + diff --git a/src/docbkx/appendix/appendix-schema.xml b/src/docbkx/appendix/appendix-schema.xml new file mode 100644 index 000000000..ef7ce9efe --- /dev/null +++ b/src/docbkx/appendix/appendix-schema.xml @@ -0,0 +1,16 @@ + + + + + Spring Data Key Value Schema(s) + + Spring Data - Redis support + + + FIXME: REDIS SCHEMA LOCATION/NAME CHANGED + + + + + diff --git a/src/docbkx/appendix/introduction.xml b/src/docbkx/appendix/introduction.xml new file mode 100644 index 000000000..6d9d95050 --- /dev/null +++ b/src/docbkx/appendix/introduction.xml @@ -0,0 +1,10 @@ + + Document structure + + + Various appendixes outside the reference documentation. + + + defines the schemas provided by Spring Data + Key Value. + \ No newline at end of file diff --git a/src/docbkx/index.xml b/src/docbkx/index.xml index a6fbafe92..d3ec64c7d 100644 --- a/src/docbkx/index.xml +++ b/src/docbkx/index.xml @@ -50,6 +50,13 @@ + + Appendixes + + + + +