This is an automated email from the ASF dual-hosted git repository.

SvenO3 pushed a commit to branch fix-adapter-producer-lifecycle
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to refs/heads/fix-adapter-producer-lifecycle 
by this push:
     new fac716e7a4 Use ConcurrentHashMaps for registries
fac716e7a4 is described below

commit fac716e7a44b51d72126715b1e64c7ebf456688c
Author: Sven Oehler <[email protected]>
AuthorDate: Mon Jul 13 17:01:45 2026 +0200

    Use ConcurrentHashMaps for registries
---
 .../management/init/RunningAdapterInstance.java    | 19 ++++++++++++-
 .../management/init/RunningAdapterInstances.java   |  9 +++---
 .../management/init/RunningInstances.java          |  4 +--
 .../standalone/manager/ProtocolManager.java        | 32 ++++++++++------------
 4 files changed, 40 insertions(+), 24 deletions(-)

diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstance.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstance.java
index 8a0a394863..02094422b3 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstance.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstance.java
@@ -1,3 +1,21 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You 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.apache.streampipes.extensions.management.init;
 
 import org.apache.streampipes.extensions.api.connect.StreamPipesAdapter;
@@ -6,4 +24,3 @@ import 
org.apache.streampipes.extensions.management.connect.adapter.model.EventC
 public record RunningAdapterInstance(StreamPipesAdapter adapter,
                                      EventCollector eventCollector) {
 }
-
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstances.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstances.java
index 959bed6777..c0229c355e 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstances.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningAdapterInstances.java
@@ -23,15 +23,15 @@ import 
org.apache.streampipes.extensions.management.connect.adapter.model.EventC
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 
 import java.util.Collection;
-import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
 
 public enum RunningAdapterInstances {
   INSTANCE;
 
-  private final Map<String, StreamPipesAdapter> runningAdapterInstances = new 
HashMap<>();
-  private final Map<String, AdapterDescription> 
runningAdapterDescriptionInstances = new HashMap<>();
-  private final Map<String, EventCollector> runningAdapterCollectors = new 
HashMap<>();
+  private final Map<String, StreamPipesAdapter> runningAdapterInstances = new 
ConcurrentHashMap<>();
+  private final Map<String, AdapterDescription> 
runningAdapterDescriptionInstances = new ConcurrentHashMap<>();
+  private final Map<String, EventCollector> runningAdapterCollectors = new 
ConcurrentHashMap<>();
 
   public void addAdapter(String elementId,
                          StreamPipesAdapter adapter,
@@ -56,4 +56,5 @@ public enum RunningAdapterInstances {
   public Collection<AdapterDescription> getAllRunningAdapterDescriptions() {
     return this.runningAdapterDescriptionInstances.values();
   }
+
 }
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningInstances.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningInstances.java
index f14a94de00..712d1de23e 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningInstances.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/init/RunningInstances.java
@@ -23,17 +23,17 @@ import 
org.apache.streampipes.extensions.management.util.ElementInfo;
 import org.apache.streampipes.model.base.NamedStreamPipesEntity;
 
 import java.util.ArrayList;
-import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.stream.Collectors;
 
 public enum RunningInstances {
   INSTANCE;
 
   private final Map<String,
-      ElementInfo<NamedStreamPipesEntity, IStreamPipesRuntime<?, ?>>> 
runningInstances = new HashMap<>();
+      ElementInfo<NamedStreamPipesEntity, IStreamPipesRuntime<?, ?>>> 
runningInstances = new ConcurrentHashMap<>();
 
 
   public void add(String id,
diff --git 
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/manager/ProtocolManager.java
 
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/manager/ProtocolManager.java
index bf86a4bfaf..3bfc3e9e4a 100644
--- 
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/manager/ProtocolManager.java
+++ 
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/manager/ProtocolManager.java
@@ -27,14 +27,14 @@ import 
org.apache.streampipes.wrapper.standalone.routing.StandaloneSpOutputColle
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
 
 public class ProtocolManager {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(ProtocolManager.class);
-  public static Map<String, StandaloneSpInputCollector> consumers = new 
HashMap<>();
-  public static Map<String, StandaloneSpOutputCollector> producers = new 
HashMap<>();
+  public static Map<String, StandaloneSpInputCollector> consumers = new 
ConcurrentHashMap<>();
+  public static Map<String, StandaloneSpOutputCollector> producers = new 
ConcurrentHashMap<>();
 
   // TODO currently only the topic name is used as an identifier for a 
consumer/producer. Should
   // be changed by some hashCode implementation in streampipes-model, but this 
requires changes
@@ -45,13 +45,12 @@ public class ProtocolManager {
       throws SpRuntimeException {
     ProtocolOverrides.addNatsTokenIfConfigured(protocol);
 
-    if (consumers.containsKey(topicName(protocol))) {
-      return consumers.get(topicName(protocol));
-    } else {
-      consumers.put(topicName(protocol), makeInputCollector(protocol, 
singletonEngine));
-      LOG.debug("Adding new consumer to consumer map (size=" + 
consumers.size() + "): " + topicName(protocol));
-      return consumers.get(topicName(protocol));
-    }
+    var topic = topicName(protocol);
+    return consumers.computeIfAbsent(topic, key -> {
+      var inputCollector = makeInputCollector(protocol, singletonEngine);
+      LOG.debug("Adding new consumer to consumer map (size={}): {}", 
consumers.size(), key);
+      return inputCollector;
+    });
 
   }
 
@@ -60,15 +59,14 @@ public class ProtocolManager {
       throws SpRuntimeException {
     ProtocolOverrides.addNatsTokenIfConfigured(protocol);
 
-    if (producers.containsKey(topicName(protocol))) {
-      return producers.get(topicName(protocol));
-    } else {
-      producers.put(topicName(protocol), makeOutputCollector(protocol, 
resourceId));
+    var topic = topicName(protocol);
+    return producers.computeIfAbsent(topic, key -> {
+      var outputCollector = makeOutputCollector(protocol, resourceId);
       LOG.debug("Adding new producer to producer map (size={}): {}",
           producers.size(),
-          topicName(protocol));
-      return producers.get(topicName(protocol));
-    }
+          key);
+      return outputCollector;
+    });
 
   }
 

Reply via email to