From f8c699beb5ed5e11d0c53e48a52e3db236c41c65 Mon Sep 17 00:00:00 2001 From: John Blum Date: Mon, 23 Oct 2023 12:18:57 -0700 Subject: [PATCH] Safely add and register the MessageListener to Topic mapping. Given addListener(:MessageListener, :Collection) could be called concurrently from the addMessageListener(:MessageListener, Collection) method by multiple Threads, and the RedisMessageListenerContainer Javadoc specifically states that it is safe to call the addMessageListener(..) method conurrently without any external synchronization, and the registeration (or mapping) of listener to Topics is a componund action, then a race condition is possible. Closes #2755 --- .../listener/RedisMessageListenerContainer.java | 13 +++++-------- 1 file changed, 5 insertions(+), 8 deletions(-) diff --git a/src/main/java/org/springframework/data/redis/listener/RedisMessageListenerContainer.java b/src/main/java/org/springframework/data/redis/listener/RedisMessageListenerContainer.java index 2b6fe773a..717f1acf6 100644 --- a/src/main/java/org/springframework/data/redis/listener/RedisMessageListenerContainer.java +++ b/src/main/java/org/springframework/data/redis/listener/RedisMessageListenerContainer.java @@ -612,20 +612,17 @@ public class RedisMessageListenerContainer implements InitializingBean, Disposab private void addListener(MessageListener listener, Collection topics) { - Assert.notNull(listener, "a valid listener is required"); - Assert.notEmpty(topics, "at least one topic is required"); + Assert.notNull(listener, "A valid listener is required"); + Assert.notEmpty(topics, "At least one topic is required"); List channels = new ArrayList<>(topics.size()); List patterns = new ArrayList<>(topics.size()); boolean trace = logger.isTraceEnabled(); - // add listener mapping - Set set = listenerTopics.get(listener); - if (set == null) { - set = new CopyOnWriteArraySet<>(); - listenerTopics.put(listener, set); - } + // safely lookup or add MessageListener to Topic mapping + Set set = listenerTopics.computeIfAbsent(listener, key -> new CopyOnWriteArraySet<>()); + set.addAll(topics); for (Topic topic : topics) {