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;
+ });
}