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

riemer pushed a commit to branch 
1717-support-other-protocols-besides-kafka-in-streampipes-client-for-gathering-live-data
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to 
refs/heads/1717-support-other-protocols-besides-kafka-in-streampipes-client-for-gathering-live-data
 by this push:
     new d74046e44 Improve handling of messaging protocols, extend broker 
support in client (#1717)
d74046e44 is described below

commit d74046e4419a9686fdb26230afe42c83316b7db8
Author: Dominik Riemer <[email protected]>
AuthorDate: Tue Jun 27 14:42:43 2023 +0200

    Improve handling of messaging protocols, extend broker support in client 
(#1717)
---
 .../streampipes/client/api/IDataProcessorApi.java  |  29 +++--
 .../streampipes/client/api/IDataSinkApi.java       |  15 ++-
 .../streampipes/client/api/IDataStreamApi.java     |  18 ++--
 .../streampipes/client/api/IStreamPipesClient.java |   3 +
 .../api/config/IStreamPipesClientConfig.java       |   5 +-
 .../client/api/live/IBrokerConfigOverride.java     |  11 +-
 ...kaConfig.java => IConfiguredEventProducer.java} |  13 ++-
 .../live/{IKafkaConfig.java => ISubscription.java} |   5 +-
 streampipes-client/pom.xml                         |   2 +-
 .../streampipes/client/StreamPipesClient.java      |   6 ++
 .../streampipes/client/api/DataProcessorApi.java   |  54 +++++-----
 .../apache/streampipes/client/api/DataSinkApi.java |  26 +++--
 .../streampipes/client/api/DataStreamApi.java      |  28 +++--
 .../client/live/ConfiguredEventProducer.java       |  53 ++++++++++
 .../streampipes/client/live/ProducerManager.java   |  62 +++++++++++
 .../streampipes/client/live/Subscription.java      |  25 +++--
 .../client/live/SubscriptionManager.java           | 117 ++++++++++-----------
 .../client/model/StreamPipesClientConfig.java      |  18 ++--
 .../elements/SendToBrokerAdapterSink.java          |  21 ++--
 .../elements/SendToJmsAdapterSink.java             |   8 +-
 .../elements/SendToKafkaAdapterSink.java           |   8 +-
 .../elements/SendToMqttAdapterSink.java            |   8 +-
 .../elements/SendToNatsAdapterSink.java            |   8 +-
 .../connect/iiot/protocol/stream/NatsProtocol.java |  14 +--
 .../sinks/brokers/jvm/jms/JmsPublisherSink.java    |   6 +-
 streampipes-integration-tests/pom.xml              |  11 ++
 .../integration/adapters/MqttAdapterTester.java    |  18 +---
 .../integration/adapters/PulsarAdapterTester.java  |  14 +--
 .../integration/client/ClientLiveDataTest.java     |  15 +--
 .../client/ClientLiveDataTesterBase.java           | 117 +++++++++++++++++++++
 .../integration/client/ClientNatsTester.java       |  64 +++++++++++
 .../integration/containers/NatsContainer.java      |  36 ++++---
 .../integration/containers/NatsDevContainer.java   |  10 +-
 .../streampipes/integration/utils/Utils.java       |  25 +++--
 .../messaging/jms/ActiveMQConnectionProvider.java  |   8 ++
 .../messaging/jms/ActiveMQConsumer.java            |  12 ++-
 .../messaging/jms/ActiveMQPublisher.java           |  47 ++-------
 .../streampipes/messaging/jms/SpJmsProtocol.java   |  15 +--
 .../messaging/kafka/SpKafkaConsumer.java           |  14 ++-
 .../messaging/kafka/SpKafkaProducer.java           |   8 +-
 .../messaging/kafka/SpKafkaProtocol.java           |  16 +--
 .../messaging/mqtt/AbstractMqttConnector.java      |   6 ++
 .../streampipes/messaging/mqtt/MqttConsumer.java   |  12 ++-
 .../streampipes/messaging/mqtt/MqttPublisher.java  |  12 ++-
 .../streampipes/messaging/mqtt/SpMqttProtocol.java |  16 +--
 .../messaging/nats/AbstractNatsConnector.java      |   8 ++
 .../streampipes/messaging/nats/NatsConsumer.java   |  22 ++--
 .../streampipes/messaging/nats/NatsPublisher.java  |  12 ++-
 .../streampipes/messaging/nats/SpNatsProtocol.java |  16 +--
 .../streampipes/messaging/EventConsumer.java       |   6 +-
 .../streampipes/messaging/EventProducer.java       |   5 +-
 .../messaging/SpProtocolDefinition.java            |   5 +-
 .../runtime/PipelineElementRuntimeInfoFetcher.java |  43 +++-----
 .../routing/StandaloneSpInputCollector.java        |  14 ++-
 .../routing/StandaloneSpOutputCollector.java       |  12 +--
 55 files changed, 746 insertions(+), 436 deletions(-)

diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataProcessorApi.java
 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataProcessorApi.java
index 8f6f38945..b0eabf750 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataProcessorApi.java
+++ 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataProcessorApi.java
@@ -21,10 +21,9 @@ package org.apache.streampipes.client.api;
 import org.apache.streampipes.client.api.annotation.NotYetImplemented;
 import org.apache.streampipes.client.api.constants.InputStreamIndex;
 import org.apache.streampipes.client.api.live.EventProcessor;
-import org.apache.streampipes.client.api.live.IKafkaConfig;
-import org.apache.streampipes.messaging.EventConsumer;
+import org.apache.streampipes.client.api.live.IBrokerConfigOverride;
+import org.apache.streampipes.client.api.live.ISubscription;
 import org.apache.streampipes.model.graph.DataProcessorInvocation;
-import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
 
 import java.util.List;
 import java.util.Optional;
@@ -49,19 +48,19 @@ public interface IDataProcessorApi extends CRUDApi<String, 
DataProcessorInvocati
   @NotYetImplemented
   void update(DataProcessorInvocation element);
 
-  EventConsumer<KafkaTransportProtocol> subscribe(DataProcessorInvocation 
processor,
-                            EventProcessor callback);
+  ISubscription subscribe(DataProcessorInvocation processor,
+                          EventProcessor callback);
 
-  EventConsumer<KafkaTransportProtocol> subscribe(DataProcessorInvocation 
processor,
-                            IKafkaConfig kafkaConfig,
-                            EventProcessor callback);
+  ISubscription subscribe(DataProcessorInvocation processor,
+                          IBrokerConfigOverride brokerConfigOverride,
+                          EventProcessor callback);
 
-  EventConsumer<KafkaTransportProtocol> subscribe(DataProcessorInvocation 
processor,
-                                                  InputStreamIndex index,
-                                                  EventProcessor callback);
+  ISubscription subscribe(DataProcessorInvocation processor,
+                          InputStreamIndex index,
+                          EventProcessor callback);
 
-  EventConsumer<KafkaTransportProtocol> subscribe(DataProcessorInvocation 
processor,
-                            InputStreamIndex index,
-                            IKafkaConfig kafkaConfig,
-                            EventProcessor callback);
+  ISubscription subscribe(DataProcessorInvocation processor,
+                          InputStreamIndex index,
+                          IBrokerConfigOverride brokerConfigOverride,
+                          EventProcessor callback);
 }
diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataSinkApi.java
 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataSinkApi.java
index f868f9cd5..b34d93945 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataSinkApi.java
+++ 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataSinkApi.java
@@ -20,10 +20,9 @@ package org.apache.streampipes.client.api;
 
 import org.apache.streampipes.client.api.annotation.NotYetImplemented;
 import org.apache.streampipes.client.api.live.EventProcessor;
-import org.apache.streampipes.client.api.live.IKafkaConfig;
-import org.apache.streampipes.messaging.EventConsumer;
+import org.apache.streampipes.client.api.live.IBrokerConfigOverride;
+import org.apache.streampipes.client.api.live.ISubscription;
 import org.apache.streampipes.model.graph.DataSinkInvocation;
-import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
 
 import java.util.List;
 import java.util.Optional;
@@ -46,12 +45,12 @@ public interface IDataSinkApi extends CRUDApi<String, 
DataSinkInvocation> {
   @Override
   void update(DataSinkInvocation element);
 
-  EventConsumer<KafkaTransportProtocol> subscribe(DataSinkInvocation sink,
-                                                  EventProcessor callback);
+  ISubscription subscribe(DataSinkInvocation sink,
+                          EventProcessor callback);
 
-  EventConsumer<KafkaTransportProtocol> subscribe(DataSinkInvocation sink,
-                                                  IKafkaConfig kafkaConfig,
-                                                  EventProcessor callback);
+  ISubscription subscribe(DataSinkInvocation sink,
+                          IBrokerConfigOverride brokerConfigOverride,
+                          EventProcessor callback);
 
   DataSinkInvocation getDataSinkForPipelineElement(String templateId, 
DataSinkInvocation pipelineElement);
 }
diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataStreamApi.java
 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataStreamApi.java
index e3fa253f2..8435511f6 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataStreamApi.java
+++ 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IDataStreamApi.java
@@ -19,10 +19,10 @@
 package org.apache.streampipes.client.api;
 
 import org.apache.streampipes.client.api.live.EventProcessor;
-import org.apache.streampipes.client.api.live.IKafkaConfig;
-import org.apache.streampipes.messaging.EventConsumer;
+import org.apache.streampipes.client.api.live.IBrokerConfigOverride;
+import org.apache.streampipes.client.api.live.IConfiguredEventProducer;
+import org.apache.streampipes.client.api.live.ISubscription;
 import org.apache.streampipes.model.SpDataStream;
-import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
 
 import java.util.List;
 import java.util.Optional;
@@ -43,10 +43,12 @@ public interface IDataStreamApi extends CRUDApi<String, 
SpDataStream> {
   @Override
   void update(SpDataStream element);
 
-  EventConsumer<KafkaTransportProtocol> subscribe(SpDataStream stream,
-                            EventProcessor callback);
+  IConfiguredEventProducer getProducer(SpDataStream stream);
 
-  EventConsumer<KafkaTransportProtocol> subscribe(SpDataStream stream,
-                                                  IKafkaConfig kafkaConfig,
-                                                  EventProcessor callback);
+  ISubscription subscribe(SpDataStream stream,
+                          EventProcessor callback);
+
+  ISubscription subscribe(SpDataStream stream,
+                          IBrokerConfigOverride brokerConfigOverride,
+                          EventProcessor callback);
 }
diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IStreamPipesClient.java
 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IStreamPipesClient.java
index 7747b28c0..b326b362b 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IStreamPipesClient.java
+++ 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/IStreamPipesClient.java
@@ -22,6 +22,7 @@ import 
org.apache.streampipes.client.api.config.ClientConnectionUrlResolver;
 import org.apache.streampipes.client.api.config.IStreamPipesClientConfig;
 import org.apache.streampipes.client.api.credentials.CredentialsProvider;
 import org.apache.streampipes.dataformat.SpDataFormatFactory;
+import org.apache.streampipes.messaging.SpProtocolDefinitionFactory;
 import org.apache.streampipes.model.mail.SpEmail;
 
 import java.io.Serializable;
@@ -29,6 +30,8 @@ import java.io.Serializable;
 public interface IStreamPipesClient extends Serializable {
   void registerDataFormat(SpDataFormatFactory spDataFormatFactory);
 
+  void registerProtocol(SpProtocolDefinitionFactory<?> 
spProtocolDefinitionFactory);
+
   CredentialsProvider getCredentials();
 
   IStreamPipesClientConfig getConfig();
diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
index b1576ce67..4238eec3d 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
+++ 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
@@ -19,17 +19,16 @@
 package org.apache.streampipes.client.api.config;
 
 import org.apache.streampipes.dataformat.SpDataFormatFactory;
+import org.apache.streampipes.messaging.SpProtocolDefinitionFactory;
 
 import com.fasterxml.jackson.databind.ObjectMapper;
 
-import java.util.List;
-
 public interface IStreamPipesClientConfig {
   ObjectMapper getSerializer();
 
   void addDataFormat(SpDataFormatFactory spDataFormatFactory);
 
-  List<SpDataFormatFactory> getRegisteredDataFormats();
+  void addTransportProtocol(SpProtocolDefinitionFactory<?> 
protocolDefinitionFactory);
 
   ClientConnectionUrlResolver getConnectionConfig();
 }
diff --git 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IBrokerConfigOverride.java
similarity index 72%
copy from 
streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
copy to 
streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IBrokerConfigOverride.java
index cd24baad5..5c68a67a1 100644
--- 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
+++ 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IBrokerConfigOverride.java
@@ -16,13 +16,16 @@
  *
  */
 
-package org.apache.streampipes.messaging;
+package org.apache.streampipes.client.api.live;
 
+import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
 import org.apache.streampipes.model.grounding.TransportProtocol;
 
-public interface SpProtocolDefinition<T extends TransportProtocol> {
+public interface IBrokerConfigOverride {
 
-  EventConsumer<T> getConsumer();
+  void overrideHostname(TransportProtocol protocol);
 
-  EventProducer<T> getProducer();
+  void overridePort(TransportProtocol protocol);
+
+  void overrideKafkaHostname(KafkaTransportProtocol protocol);
 }
diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IKafkaConfig.java
 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IConfiguredEventProducer.java
similarity index 80%
copy from 
streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IKafkaConfig.java
copy to 
streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IConfiguredEventProducer.java
index 6e70c06fe..1be40fc0c 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IKafkaConfig.java
+++ 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IConfiguredEventProducer.java
@@ -18,8 +18,15 @@
 
 package org.apache.streampipes.client.api.live;
 
-public interface IKafkaConfig {
-  String getKafkaHost();
+import org.apache.streampipes.model.runtime.Event;
 
-  Integer getKafkaPort();
+import java.util.Map;
+
+public interface IConfiguredEventProducer {
+
+  void publish(Event event);
+
+  void publish(Map<String, Object> event);
+
+  void close();
 }
diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IKafkaConfig.java
 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/ISubscription.java
similarity index 91%
copy from 
streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IKafkaConfig.java
copy to 
streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/ISubscription.java
index 6e70c06fe..cbd173b42 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IKafkaConfig.java
+++ 
b/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/ISubscription.java
@@ -18,8 +18,7 @@
 
 package org.apache.streampipes.client.api.live;
 
-public interface IKafkaConfig {
-  String getKafkaHost();
+public interface ISubscription {
 
-  Integer getKafkaPort();
+  void unsubscribe();
 }
diff --git a/streampipes-client/pom.xml b/streampipes-client/pom.xml
index 263ffa43d..3c47a7edc 100644
--- a/streampipes-client/pom.xml
+++ b/streampipes-client/pom.xml
@@ -51,7 +51,7 @@
         </dependency>
         <dependency>
             <groupId>org.apache.streampipes</groupId>
-            <artifactId>streampipes-messaging-kafka</artifactId>
+            <artifactId>streampipes-messaging</artifactId>
             <version>0.93.0-SNAPSHOT</version>
         </dependency>
         <dependency>
diff --git 
a/streampipes-client/src/main/java/org/apache/streampipes/client/StreamPipesClient.java
 
b/streampipes-client/src/main/java/org/apache/streampipes/client/StreamPipesClient.java
index b6583e60e..1eee9f475 100644
--- 
a/streampipes-client/src/main/java/org/apache/streampipes/client/StreamPipesClient.java
+++ 
b/streampipes-client/src/main/java/org/apache/streampipes/client/StreamPipesClient.java
@@ -40,6 +40,7 @@ import org.apache.streampipes.dataformat.SpDataFormatFactory;
 import org.apache.streampipes.dataformat.cbor.CborDataFormatFactory;
 import org.apache.streampipes.dataformat.fst.FstDataFormatFactory;
 import org.apache.streampipes.dataformat.json.JsonDataFormatFactory;
+import org.apache.streampipes.messaging.SpProtocolDefinitionFactory;
 import org.apache.streampipes.model.mail.SpEmail;
 
 public class StreamPipesClient implements
@@ -134,6 +135,11 @@ public class StreamPipesClient implements
     this.config.addDataFormat(spDataFormatFactory);
   }
 
+  @Override
+  public void registerProtocol(SpProtocolDefinitionFactory<?> 
spProtocolDefinitionFactory) {
+    this.config.addTransportProtocol(spProtocolDefinitionFactory);
+  }
+
   @Override
   public CredentialsProvider getCredentials() {
     return config.getConnectionConfig().getCredentials();
diff --git 
a/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataProcessorApi.java
 
b/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataProcessorApi.java
index 099070e28..f3c3e1722 100644
--- 
a/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataProcessorApi.java
+++ 
b/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataProcessorApi.java
@@ -20,11 +20,11 @@ package org.apache.streampipes.client.api;
 import org.apache.streampipes.client.api.annotation.NotYetImplemented;
 import org.apache.streampipes.client.api.constants.InputStreamIndex;
 import org.apache.streampipes.client.api.live.EventProcessor;
-import org.apache.streampipes.client.api.live.IKafkaConfig;
+import org.apache.streampipes.client.api.live.IBrokerConfigOverride;
+import org.apache.streampipes.client.api.live.ISubscription;
 import org.apache.streampipes.client.live.SubscriptionManager;
 import org.apache.streampipes.client.model.StreamPipesClientConfig;
 import org.apache.streampipes.client.util.StreamPipesApiPath;
-import org.apache.streampipes.messaging.kafka.SpKafkaConsumer;
 import org.apache.streampipes.model.graph.DataProcessorInvocation;
 
 import java.util.List;
@@ -77,23 +77,23 @@ public class DataProcessorApi extends 
AbstractTypedClientApi<DataProcessorInvoca
    * @param callback  The callback where events will be received
    */
   @Override
-  public SpKafkaConsumer subscribe(DataProcessorInvocation processor,
-                                   EventProcessor callback) {
-    return new SubscriptionManager(clientConfig, 
processor.getOutputStream().getEventGrounding(), callback).subscribe();
+  public ISubscription subscribe(DataProcessorInvocation processor,
+                                 EventProcessor callback) {
+    return new 
SubscriptionManager(processor.getOutputStream().getEventGrounding(), 
callback).subscribe();
   }
 
   /**
    * Subscribe to the output stream of the processor
    *
-   * @param processor   The data processor to subscribe to
-   * @param kafkaConfig Additional kafka settings which will override the 
default value (see docs)
-   * @param callback    The callback where events will be received
+   * @param processor            The data processor to subscribe to
+   * @param brokerConfigOverride Additional broker settings which will 
override the default value (see docs)
+   * @param callback             The callback where events will be received
    */
   @Override
-  public SpKafkaConsumer subscribe(DataProcessorInvocation processor,
-                                   IKafkaConfig kafkaConfig,
-                                   EventProcessor callback) {
-    return new SubscriptionManager(clientConfig, kafkaConfig, 
processor.getOutputStream().getEventGrounding(), callback)
+  public ISubscription subscribe(DataProcessorInvocation processor,
+                                 IBrokerConfigOverride brokerConfigOverride,
+                                 EventProcessor callback) {
+    return new SubscriptionManager(brokerConfigOverride, 
processor.getOutputStream().getEventGrounding(), callback)
         .subscribe();
   }
 
@@ -105,27 +105,29 @@ public class DataProcessorApi extends 
AbstractTypedClientApi<DataProcessorInvoca
    * @param callback  The callback where events will be received
    */
   @Override
-  public SpKafkaConsumer subscribe(DataProcessorInvocation processor,
-                                   InputStreamIndex index,
-                                   EventProcessor callback) {
-    return new SubscriptionManager(clientConfig,
-        processor.getInputStreams().get(index.toIndex()).getEventGrounding(), 
callback).subscribe();
+  public ISubscription subscribe(DataProcessorInvocation processor,
+                                 InputStreamIndex index,
+                                 EventProcessor callback) {
+    return new SubscriptionManager(
+        processor.getInputStreams().get(index.toIndex()).getEventGrounding(), 
callback)
+        .subscribe();
   }
 
   /**
    * Subscribe to the input stream of the sink
    *
-   * @param processor   The data processor to subscribe to
-   * @param index       The index of the input stream
-   * @param kafkaConfig Additional kafka settings which will override the 
default value (see docs)
-   * @param callback    The callback where events will be received
+   * @param processor            The data processor to subscribe to
+   * @param index                The index of the input stream
+   * @param brokerConfigOverride Additional kafka settings which will override 
the default value (see docs)
+   * @param callback             The callback where events will be received
    */
   @Override
-  public SpKafkaConsumer subscribe(DataProcessorInvocation processor,
-                                   InputStreamIndex index,
-                                   IKafkaConfig kafkaConfig,
-                                   EventProcessor callback) {
-    return new SubscriptionManager(clientConfig, kafkaConfig,
+  public ISubscription subscribe(DataProcessorInvocation processor,
+                                 InputStreamIndex index,
+                                 IBrokerConfigOverride brokerConfigOverride,
+                                 EventProcessor callback) {
+    return new SubscriptionManager(
+        brokerConfigOverride,
         processor.getInputStreams().get(index.toIndex()).getEventGrounding(), 
callback)
         .subscribe();
   }
diff --git 
a/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataSinkApi.java
 
b/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataSinkApi.java
index 48200a2ca..e43c4eace 100644
--- 
a/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataSinkApi.java
+++ 
b/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataSinkApi.java
@@ -19,13 +19,12 @@ package org.apache.streampipes.client.api;
 
 import org.apache.streampipes.client.api.annotation.NotYetImplemented;
 import org.apache.streampipes.client.api.live.EventProcessor;
-import org.apache.streampipes.client.api.live.IKafkaConfig;
+import org.apache.streampipes.client.api.live.IBrokerConfigOverride;
+import org.apache.streampipes.client.api.live.ISubscription;
 import org.apache.streampipes.client.live.SubscriptionManager;
 import org.apache.streampipes.client.model.StreamPipesClientConfig;
 import org.apache.streampipes.client.util.StreamPipesApiPath;
-import org.apache.streampipes.messaging.EventConsumer;
 import org.apache.streampipes.model.graph.DataSinkInvocation;
-import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
 
 import java.util.List;
 import java.util.Optional;
@@ -71,24 +70,23 @@ public class DataSinkApi extends 
AbstractTypedClientApi<DataSinkInvocation>
    * @param callback The callback where events will be received
    */
   @Override
-  public EventConsumer<KafkaTransportProtocol> subscribe(DataSinkInvocation 
sink,
-                                                         EventProcessor 
callback) {
-    return new SubscriptionManager(clientConfig,
-        sink.getInputStreams().get(0).getEventGrounding(), 
callback).subscribe();
+  public ISubscription subscribe(DataSinkInvocation sink,
+                                 EventProcessor callback) {
+    return new 
SubscriptionManager(sink.getInputStreams().get(0).getEventGrounding(), 
callback).subscribe();
   }
 
   /**
    * Subscribe to the input stream of the sink
    *
-   * @param sink        The data sink to subscribe to
-   * @param kafkaConfig Additional kafka settings which will override the 
default value (see docs)
-   * @param callback    The callback where events will be received
+   * @param sink                 The data sink to subscribe to
+   * @param brokerConfigOverride Additional kafka settings which will override 
the default value (see docs)
+   * @param callback             The callback where events will be received
    */
   @Override
-  public EventConsumer<KafkaTransportProtocol> subscribe(DataSinkInvocation 
sink,
-                                   IKafkaConfig kafkaConfig,
-                                   EventProcessor callback) {
-    return new SubscriptionManager(clientConfig, kafkaConfig,
+  public ISubscription subscribe(DataSinkInvocation sink,
+                                 IBrokerConfigOverride brokerConfigOverride,
+                                 EventProcessor callback) {
+    return new SubscriptionManager(brokerConfigOverride,
         sink.getInputStreams().get(0).getEventGrounding(), 
callback).subscribe();
   }
 
diff --git 
a/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataStreamApi.java
 
b/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataStreamApi.java
index d49e0a0c7..950f53a9a 100644
--- 
a/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataStreamApi.java
+++ 
b/streampipes-client/src/main/java/org/apache/streampipes/client/api/DataStreamApi.java
@@ -18,13 +18,14 @@
 package org.apache.streampipes.client.api;
 
 import org.apache.streampipes.client.api.live.EventProcessor;
-import org.apache.streampipes.client.api.live.IKafkaConfig;
+import org.apache.streampipes.client.api.live.IBrokerConfigOverride;
+import org.apache.streampipes.client.api.live.IConfiguredEventProducer;
+import org.apache.streampipes.client.api.live.ISubscription;
+import org.apache.streampipes.client.live.ProducerManager;
 import org.apache.streampipes.client.live.SubscriptionManager;
 import org.apache.streampipes.client.model.StreamPipesClientConfig;
 import org.apache.streampipes.client.util.StreamPipesApiPath;
-import org.apache.streampipes.messaging.EventConsumer;
 import org.apache.streampipes.model.SpDataStream;
-import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
 import org.apache.streampipes.model.message.Message;
 
 import java.net.URLEncoder;
@@ -56,7 +57,7 @@ public class DataStreamApi extends 
AbstractTypedClientApi<SpDataStream> implemen
   /**
    * Directly install a new data stream
    *
-   * @param stream   The data stream to add
+   * @param stream The data stream to add
    */
   @Override
   public void create(SpDataStream stream) {
@@ -78,6 +79,11 @@ public class DataStreamApi extends 
AbstractTypedClientApi<SpDataStream> implemen
 
   }
 
+  @Override
+  public IConfiguredEventProducer getProducer(SpDataStream stream) {
+    return new ProducerManager(stream.getEventGrounding()).makeProducer();
+  }
+
   /**
    * Subscribe to a data stream
    *
@@ -85,9 +91,9 @@ public class DataStreamApi extends 
AbstractTypedClientApi<SpDataStream> implemen
    * @param callback The callback where events will be received
    */
   @Override
-  public EventConsumer<KafkaTransportProtocol> subscribe(SpDataStream stream,
-                                   EventProcessor callback) {
-    return new SubscriptionManager(clientConfig, stream.getEventGrounding(), 
callback).subscribe();
+  public ISubscription subscribe(SpDataStream stream,
+                                 EventProcessor callback) {
+    return new SubscriptionManager(stream.getEventGrounding(), 
callback).subscribe();
   }
 
   /**
@@ -98,10 +104,10 @@ public class DataStreamApi extends 
AbstractTypedClientApi<SpDataStream> implemen
    * @param callback    The callback where events will be received
    */
   @Override
-  public EventConsumer<KafkaTransportProtocol> subscribe(SpDataStream stream,
-                                                         IKafkaConfig 
kafkaConfig,
-                                                         EventProcessor 
callback) {
-    return new SubscriptionManager(clientConfig, kafkaConfig, 
stream.getEventGrounding(), callback).subscribe();
+  public ISubscription subscribe(SpDataStream stream,
+                                 IBrokerConfigOverride kafkaConfig,
+                                 EventProcessor callback) {
+    return new SubscriptionManager(kafkaConfig, stream.getEventGrounding(), 
callback).subscribe();
   }
 
   @Override
diff --git 
a/streampipes-client/src/main/java/org/apache/streampipes/client/live/ConfiguredEventProducer.java
 
b/streampipes-client/src/main/java/org/apache/streampipes/client/live/ConfiguredEventProducer.java
new file mode 100644
index 000000000..d45b02e6b
--- /dev/null
+++ 
b/streampipes-client/src/main/java/org/apache/streampipes/client/live/ConfiguredEventProducer.java
@@ -0,0 +1,53 @@
+/*
+ * 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.client.live;
+
+import org.apache.streampipes.client.api.live.IConfiguredEventProducer;
+import org.apache.streampipes.dataformat.SpDataFormatDefinition;
+import org.apache.streampipes.messaging.EventProducer;
+import org.apache.streampipes.model.runtime.Event;
+
+import java.util.Map;
+
+public class ConfiguredEventProducer implements IConfiguredEventProducer {
+
+  private final EventProducer internalProducer;
+  private final SpDataFormatDefinition dataFormatDefinition;
+
+  public ConfiguredEventProducer(EventProducer internalProducer,
+                                 SpDataFormatDefinition definition) {
+    this.internalProducer = internalProducer;
+    this.dataFormatDefinition = definition;
+  }
+
+  @Override
+  public void publish(Event event) {
+    publish(event.getRaw());
+  }
+
+  @Override
+  public void publish(Map<String, Object> event) {
+    this.internalProducer.publish(dataFormatDefinition.fromMap(event));
+  }
+
+  @Override
+  public void close() {
+    this.internalProducer.disconnect();
+  }
+}
diff --git 
a/streampipes-client/src/main/java/org/apache/streampipes/client/live/ProducerManager.java
 
b/streampipes-client/src/main/java/org/apache/streampipes/client/live/ProducerManager.java
new file mode 100644
index 000000000..6cf97026d
--- /dev/null
+++ 
b/streampipes-client/src/main/java/org/apache/streampipes/client/live/ProducerManager.java
@@ -0,0 +1,62 @@
+/*
+ * 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.client.live;
+
+import org.apache.streampipes.client.api.live.IConfiguredEventProducer;
+import org.apache.streampipes.dataformat.SpDataFormatDefinition;
+import org.apache.streampipes.dataformat.SpDataFormatManager;
+import org.apache.streampipes.messaging.EventProducer;
+import org.apache.streampipes.messaging.SpProtocolManager;
+import org.apache.streampipes.model.grounding.EventGrounding;
+
+public class ProducerManager {
+
+  private final EventGrounding grounding;
+
+  public ProducerManager(EventGrounding grounding) {
+    this.grounding = grounding;
+  }
+
+  public IConfiguredEventProducer makeProducer() {
+    EventProducer producer = findProducer();
+    producer.connect();
+
+    return new ConfiguredEventProducer(
+        producer,
+        findFormatDefinition()
+    );
+  }
+
+  private EventProducer findProducer() {
+    var protocol = grounding.getTransportProtocol();
+    return SpProtocolManager
+        .INSTANCE
+        .findDefinition(protocol)
+        .orElseThrow()
+        .getProducer(protocol);
+  }
+
+  private SpDataFormatDefinition findFormatDefinition() {
+    var format = grounding.getTransportFormats().get(0);
+    return SpDataFormatManager
+        .INSTANCE
+        .findDefinition(format)
+        .orElseThrow();
+  }
+}
diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
 
b/streampipes-client/src/main/java/org/apache/streampipes/client/live/Subscription.java
similarity index 64%
copy from 
streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
copy to 
streampipes-client/src/main/java/org/apache/streampipes/client/live/Subscription.java
index b1576ce67..b0c2fdb94 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
+++ 
b/streampipes-client/src/main/java/org/apache/streampipes/client/live/Subscription.java
@@ -16,20 +16,23 @@
  *
  */
 
-package org.apache.streampipes.client.api.config;
+package org.apache.streampipes.client.live;
 
-import org.apache.streampipes.dataformat.SpDataFormatFactory;
+import org.apache.streampipes.client.api.live.ISubscription;
+import org.apache.streampipes.messaging.EventConsumer;
 
-import com.fasterxml.jackson.databind.ObjectMapper;
+public class Subscription implements ISubscription {
 
-import java.util.List;
+  private final EventConsumer consumer;
 
-public interface IStreamPipesClientConfig {
-  ObjectMapper getSerializer();
+  public Subscription(EventConsumer consumer) {
+    this.consumer = consumer;
+  }
 
-  void addDataFormat(SpDataFormatFactory spDataFormatFactory);
-
-  List<SpDataFormatFactory> getRegisteredDataFormats();
-
-  ClientConnectionUrlResolver getConnectionConfig();
+  @Override
+  public void unsubscribe() {
+    if (consumer.isConnected()) {
+      consumer.disconnect();
+    }
+  }
 }
diff --git 
a/streampipes-client/src/main/java/org/apache/streampipes/client/live/SubscriptionManager.java
 
b/streampipes-client/src/main/java/org/apache/streampipes/client/live/SubscriptionManager.java
index e58cabf0d..961d9d300 100644
--- 
a/streampipes-client/src/main/java/org/apache/streampipes/client/live/SubscriptionManager.java
+++ 
b/streampipes-client/src/main/java/org/apache/streampipes/client/live/SubscriptionManager.java
@@ -18,94 +18,93 @@
 package org.apache.streampipes.client.live;
 
 import org.apache.streampipes.client.api.live.EventProcessor;
-import org.apache.streampipes.client.api.live.IKafkaConfig;
-import org.apache.streampipes.client.model.StreamPipesClientConfig;
+import org.apache.streampipes.client.api.live.IBrokerConfigOverride;
+import org.apache.streampipes.client.api.live.ISubscription;
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
 import org.apache.streampipes.dataformat.SpDataFormatDefinition;
-import org.apache.streampipes.dataformat.SpDataFormatFactory;
-import org.apache.streampipes.messaging.kafka.SpKafkaConsumer;
+import org.apache.streampipes.dataformat.SpDataFormatManager;
+import org.apache.streampipes.messaging.EventConsumer;
+import org.apache.streampipes.messaging.SpProtocolDefinition;
+import org.apache.streampipes.messaging.SpProtocolManager;
 import org.apache.streampipes.model.grounding.EventGrounding;
 import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
+import org.apache.streampipes.model.grounding.TransportProtocol;
 import org.apache.streampipes.model.runtime.Event;
 import org.apache.streampipes.model.runtime.EventFactory;
 
-import java.util.Optional;
+import java.util.NoSuchElementException;
 
 public class SubscriptionManager {
 
   private final EventGrounding grounding;
   private final EventProcessor callback;
-  private final StreamPipesClientConfig clientConfig;
 
-  private IKafkaConfig kafkaConfig;
-  private boolean overrideKafkaSettings = false;
+  private IBrokerConfigOverride brokerConfigOverride;
+  private boolean overrideSettings = false;
 
-  public SubscriptionManager(StreamPipesClientConfig clientConfig,
-                             EventGrounding grounding,
+  public SubscriptionManager(EventGrounding grounding,
                              EventProcessor callback) {
     this.grounding = grounding;
     this.callback = callback;
-    this.clientConfig = clientConfig;
   }
 
-  public SubscriptionManager(StreamPipesClientConfig clientConfig,
-                             IKafkaConfig kafkaConfig,
+  public SubscriptionManager(IBrokerConfigOverride brokerConfigOverride,
                              EventGrounding grounding,
                              EventProcessor callback) {
-    this(clientConfig, grounding, callback);
-    this.kafkaConfig = kafkaConfig;
-    this.overrideKafkaSettings = true;
+    this(grounding, callback);
+    this.brokerConfigOverride = brokerConfigOverride;
+    this.overrideSettings = true;
   }
 
-  public SpKafkaConsumer subscribe() {
-    Optional<SpDataFormatFactory> formatConverterOpt = this
-        .clientConfig
-        .getRegisteredDataFormats()
-        .stream()
-        .filter(format -> this.grounding
-            .getTransportFormats()
-            .get(0)
-            .getRdfType()
-            .stream()
-            .anyMatch(tf -> 
tf.toString().equals(format.getTransportFormatRdfUri())))
-        .findFirst();
-
-    if (formatConverterOpt.isPresent()) {
-      final SpDataFormatDefinition converter = 
formatConverterOpt.get().createInstance();
-
-      KafkaTransportProtocol protocol =
-          overrideKafkaSettings ? overrideHostname(getKafkaProtocol()) : 
getKafkaProtocol();
-      SpKafkaConsumer kafkaConsumer = new SpKafkaConsumer(protocol, 
getOutputTopic(), event -> {
-        try {
-          Event spEvent = EventFactory.fromMap(converter.toMap(event));
-          callback.onEvent(spEvent);
-        } catch (SpRuntimeException e) {
-          e.printStackTrace();
+  public ISubscription subscribe() {
+    var formatDefinitionOpt = SpDataFormatManager
+        .INSTANCE
+        .findDefinition(this.grounding.getTransportFormats().get(0));
+
+    try {
+      SpProtocolDefinition<TransportProtocol> protocolDefinition = 
findProtocol(getTransportProtocol());
+
+      if (formatDefinitionOpt.isPresent()) {
+        final SpDataFormatDefinition converter = formatDefinitionOpt.get();
+
+        var protocol = getTransportProtocol();
+        if (overrideSettings) {
+          if (protocol instanceof KafkaTransportProtocol) {
+            
brokerConfigOverride.overrideKafkaHostname((KafkaTransportProtocol) protocol);
+          }
+          brokerConfigOverride.overrideHostname(protocol);
+          brokerConfigOverride.overridePort(protocol);
         }
-      });
-      Thread t = new Thread(kafkaConsumer);
-      t.start();
-      return kafkaConsumer;
-    } else {
+
+        EventConsumer consumer = protocolDefinition.getConsumer(protocol);
+        consumer.connect(event -> {
+          try {
+            Event spEvent = EventFactory.fromMap(converter.toMap(event));
+            callback.onEvent(spEvent);
+          } catch (SpRuntimeException e) {
+            e.printStackTrace();
+          }
+        });
+
+        return new Subscription(consumer);
+      } else {
+        throw new SpRuntimeException(
+            "No converter found for data format - did you add a format factory 
(client.registerDataFormat)?");
+      }
+    } catch (NoSuchElementException e) {
       throw new SpRuntimeException(
-          "No converter found for data format - did you add a format factory 
(client.registerDataFormat)?");
-    }
-  }
+          "Could not find an implementation for messaging protocol "
+              + 
this.grounding.getTransportProtocol().getClass().getCanonicalName()
+              + "- please add the corresponding module 
(streampipes-messaging-*) to your project dependencies.");
 
-  private KafkaTransportProtocol overrideHostname(KafkaTransportProtocol 
protocol) {
-    protocol.setBrokerHostname(kafkaConfig.getKafkaHost());
-    protocol.setKafkaPort(kafkaConfig.getKafkaPort());
-    return protocol;
+    }
   }
 
-  private KafkaTransportProtocol getKafkaProtocol() {
-    return (KafkaTransportProtocol) this.grounding.getTransportProtocol();
+  private SpProtocolDefinition<TransportProtocol> 
findProtocol(TransportProtocol protocol) {
+    return SpProtocolManager.INSTANCE.findDefinition(protocol).orElseThrow();
   }
 
-  private String getOutputTopic() {
-    return this.grounding
-        .getTransportProtocol()
-        .getTopicDefinition()
-        .getActualTopicName();
+  private TransportProtocol getTransportProtocol() {
+    return this.grounding.getTransportProtocol();
   }
 }
diff --git 
a/streampipes-client/src/main/java/org/apache/streampipes/client/model/StreamPipesClientConfig.java
 
b/streampipes-client/src/main/java/org/apache/streampipes/client/model/StreamPipesClientConfig.java
index aa36b5304..716ed50b2 100644
--- 
a/streampipes-client/src/main/java/org/apache/streampipes/client/model/StreamPipesClientConfig.java
+++ 
b/streampipes-client/src/main/java/org/apache/streampipes/client/model/StreamPipesClientConfig.java
@@ -20,23 +20,21 @@ package org.apache.streampipes.client.model;
 import org.apache.streampipes.client.api.config.ClientConnectionUrlResolver;
 import org.apache.streampipes.client.api.config.IStreamPipesClientConfig;
 import org.apache.streampipes.dataformat.SpDataFormatFactory;
+import org.apache.streampipes.dataformat.SpDataFormatManager;
+import org.apache.streampipes.messaging.SpProtocolDefinitionFactory;
+import org.apache.streampipes.messaging.SpProtocolManager;
 import org.apache.streampipes.serializers.json.JacksonSerializer;
 
 import com.fasterxml.jackson.databind.ObjectMapper;
 
-import java.util.ArrayList;
-import java.util.List;
-
 public class StreamPipesClientConfig implements IStreamPipesClientConfig {
 
-  private ClientConnectionUrlResolver connectionConfig;
-  private ObjectMapper serializer;
-  private List<SpDataFormatFactory> registeredDataFormats;
+  private final ClientConnectionUrlResolver connectionConfig;
+  private final ObjectMapper serializer;
 
   public StreamPipesClientConfig(ClientConnectionUrlResolver connectionConfig) 
{
     this.connectionConfig = connectionConfig;
     this.serializer = JacksonSerializer.getObjectMapper();
-    this.registeredDataFormats = new ArrayList<>();
   }
 
   @Override
@@ -46,12 +44,12 @@ public class StreamPipesClientConfig implements 
IStreamPipesClientConfig {
 
   @Override
   public void addDataFormat(SpDataFormatFactory spDataFormatFactory) {
-    this.registeredDataFormats.add(spDataFormatFactory);
+    SpDataFormatManager.INSTANCE.register(spDataFormatFactory);
   }
 
   @Override
-  public List<SpDataFormatFactory> getRegisteredDataFormats() {
-    return registeredDataFormats;
+  public void addTransportProtocol(SpProtocolDefinitionFactory<?> 
protocolDefinitionFactory) {
+    SpProtocolManager.INSTANCE.register(protocolDefinitionFactory);
   }
 
   @Override
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
index afda58a01..52f864db7 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToBrokerAdapterSink.java
@@ -32,20 +32,17 @@ import 
org.apache.streampipes.model.grounding.TransportProtocol;
 import org.apache.streampipes.model.monitoring.SpLogEntry;
 
 import java.util.Map;
-import java.util.function.Supplier;
 
 public abstract class SendToBrokerAdapterSink<T extends TransportProtocol> 
implements IAdapterPipelineElement {
 
   protected AdapterDescription adapterDescription;
   protected SpDataFormatDefinition dataFormatDefinition;
   protected T protocol;
-  private EventProducer<T> producer;
+  private final EventProducer producer;
 
   public SendToBrokerAdapterSink(AdapterDescription adapterDescription,
-                                 Supplier<EventProducer<T>> producerSupplier,
                                  Class<T> protocolClass) {
     this.adapterDescription = adapterDescription;
-    this.producer = producerSupplier.get();
     this.protocol = protocolClass.cast(adapterDescription
         .getEventGrounding()
         .getTransportProtocol());
@@ -54,6 +51,8 @@ public abstract class SendToBrokerAdapterSink<T extends 
TransportProtocol> imple
       modifyProtocolForDebugging(this.protocol);
     }
 
+    this.producer = makeProducer(this.protocol);
+
     TransportFormat transportFormat = adapterDescription
         .getEventGrounding()
         .getTransportFormats()
@@ -63,7 +62,7 @@ public abstract class SendToBrokerAdapterSink<T extends 
TransportProtocol> imple
         new TransportFormatSelector(transportFormat).getDataFormatDefinition();
 
     try {
-      producer.connect(protocol);
+      producer.connect();
     } catch (SpRuntimeException e) {
       e.printStackTrace();
     }
@@ -90,17 +89,9 @@ public abstract class SendToBrokerAdapterSink<T extends 
TransportProtocol> imple
     producer.publish(event);
   }
 
-  public abstract void modifyProtocolForDebugging(T transportProtocol);
+  protected abstract EventProducer makeProducer(T protocol);
 
-  public void changeTransportProtocol(T transportProtocol) {
-    try {
-      modifyProtocolForDebugging(transportProtocol);
-      producer.disconnect();
-      producer.connect(transportProtocol);
-    } catch (SpRuntimeException e) {
-      e.printStackTrace();
-    }
-  }
+  public abstract void modifyProtocolForDebugging(T protocol);
 
   private Environment getEnvironment() {
     return Environments.getEnvironment();
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToJmsAdapterSink.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToJmsAdapterSink.java
index 17f954769..98af0b42c 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToJmsAdapterSink.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToJmsAdapterSink.java
@@ -18,6 +18,7 @@
 package 
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements;
 
 import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement;
+import org.apache.streampipes.messaging.EventProducer;
 import org.apache.streampipes.messaging.jms.ActiveMQPublisher;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 import org.apache.streampipes.model.grounding.JmsTransportProtocol;
@@ -26,7 +27,12 @@ public class SendToJmsAdapterSink extends 
SendToBrokerAdapterSink<JmsTransportPr
     implements IAdapterPipelineElement {
 
   public SendToJmsAdapterSink(AdapterDescription adapterDescription) {
-    super(adapterDescription, ActiveMQPublisher::new, 
JmsTransportProtocol.class);
+    super(adapterDescription, JmsTransportProtocol.class);
+  }
+
+  @Override
+  protected EventProducer makeProducer(JmsTransportProtocol protocol) {
+    return new ActiveMQPublisher(protocol);
   }
 
   @Override
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToKafkaAdapterSink.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToKafkaAdapterSink.java
index a619140f5..5db6d32b2 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToKafkaAdapterSink.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToKafkaAdapterSink.java
@@ -18,6 +18,7 @@
 package 
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements;
 
 import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement;
+import org.apache.streampipes.messaging.EventProducer;
 import org.apache.streampipes.messaging.kafka.SpKafkaProducer;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
@@ -26,7 +27,12 @@ public class SendToKafkaAdapterSink extends 
SendToBrokerAdapterSink<KafkaTranspo
     implements IAdapterPipelineElement {
 
   public SendToKafkaAdapterSink(AdapterDescription adapterDescription) {
-    super(adapterDescription, SpKafkaProducer::new, 
KafkaTransportProtocol.class);
+    super(adapterDescription, KafkaTransportProtocol.class);
+  }
+
+  @Override
+  protected EventProducer makeProducer(KafkaTransportProtocol protocol) {
+    return new SpKafkaProducer(protocol);
   }
 
   @Override
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToMqttAdapterSink.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToMqttAdapterSink.java
index 04a481434..d2823c1c5 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToMqttAdapterSink.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToMqttAdapterSink.java
@@ -18,6 +18,7 @@
 package 
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements;
 
 import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement;
+import org.apache.streampipes.messaging.EventProducer;
 import org.apache.streampipes.messaging.mqtt.MqttPublisher;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 import org.apache.streampipes.model.grounding.MqttTransportProtocol;
@@ -26,7 +27,12 @@ public class SendToMqttAdapterSink extends 
SendToBrokerAdapterSink<MqttTransport
     implements IAdapterPipelineElement {
 
   public SendToMqttAdapterSink(AdapterDescription adapterDescription) {
-    super(adapterDescription, MqttPublisher::new, MqttTransportProtocol.class);
+    super(adapterDescription, MqttTransportProtocol.class);
+  }
+
+  @Override
+  protected EventProducer makeProducer(MqttTransportProtocol protocol) {
+    return new MqttPublisher(protocol);
   }
 
   @Override
diff --git 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToNatsAdapterSink.java
 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToNatsAdapterSink.java
index 932183c18..5216ee0b0 100644
--- 
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToNatsAdapterSink.java
+++ 
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToNatsAdapterSink.java
@@ -19,6 +19,7 @@
 package 
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements;
 
 import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement;
+import org.apache.streampipes.messaging.EventProducer;
 import org.apache.streampipes.messaging.nats.NatsPublisher;
 import org.apache.streampipes.model.connect.adapter.AdapterDescription;
 import org.apache.streampipes.model.grounding.NatsTransportProtocol;
@@ -27,7 +28,12 @@ public class SendToNatsAdapterSink extends 
SendToBrokerAdapterSink<NatsTransport
     implements IAdapterPipelineElement {
 
   public SendToNatsAdapterSink(AdapterDescription adapterDescription) {
-    super(adapterDescription, NatsPublisher::new, NatsTransportProtocol.class);
+    super(adapterDescription, NatsTransportProtocol.class);
+  }
+
+  @Override
+  protected EventProducer makeProducer(NatsTransportProtocol protocol) {
+    return new NatsPublisher(protocol);
   }
 
   @Override
diff --git 
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/NatsProtocol.java
 
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/NatsProtocol.java
index 52cd2f96e..ce3a2a778 100644
--- 
a/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/NatsProtocol.java
+++ 
b/streampipes-extensions/streampipes-connect-adapters-iiot/src/main/java/org/apache/streampipes/connect/iiot/protocol/stream/NatsProtocol.java
@@ -18,6 +18,7 @@
 
 package org.apache.streampipes.connect.iiot.protocol.stream;
 
+import org.apache.streampipes.commons.exceptions.SpRuntimeException;
 import org.apache.streampipes.commons.exceptions.connect.AdapterException;
 import org.apache.streampipes.commons.exceptions.connect.ParseException;
 import org.apache.streampipes.extensions.api.connect.IAdapterConfiguration;
@@ -43,7 +44,6 @@ import org.apache.streampipes.sdk.helpers.Locales;
 import org.apache.streampipes.sdk.utils.Assets;
 
 import java.io.ByteArrayInputStream;
-import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
 import java.util.concurrent.TimeUnit;
@@ -126,10 +126,10 @@ public class NatsProtocol implements StreamPipesAdapter {
                                IEventCollector collector,
                                IAdapterRuntimeContext adapterRuntimeContext) 
throws AdapterException {
     this.applyConfiguration(extractor.getStaticPropertyExtractor());
-    this.natsConsumer = new NatsConsumer();
+    this.natsConsumer = new NatsConsumer(natsConfig);
     try {
-      this.natsConsumer.connect(natsConfig, new 
BrokerEventProcessor(extractor.selectedParser(), collector));
-    } catch (IOException | InterruptedException e) {
+      this.natsConsumer.connect(new 
BrokerEventProcessor(extractor.selectedParser(), collector));
+    } catch (SpRuntimeException e) {
       throw new AdapterException("Error when connecting to the Nats broker on "
           + natsConfig.getNatsUrls() + " . ", e);
     }
@@ -146,7 +146,7 @@ public class NatsProtocol implements StreamPipesAdapter {
                                        IAdapterGuessSchemaContext 
adapterGuessSchemaContext) throws AdapterException {
     this.applyConfiguration(extractor.getStaticPropertyExtractor());
     List<byte[]> elements = new ArrayList<>();
-    this.natsConsumer = new NatsConsumer();
+    this.natsConsumer = new NatsConsumer(natsConfig);
     final boolean[] completed = {false};
     InternalEventProcessor<byte[]> processor = event -> {
       elements.add(event);
@@ -154,8 +154,8 @@ public class NatsProtocol implements StreamPipesAdapter {
     };
 
     try {
-      this.natsConsumer.connect(natsConfig, processor);
-    } catch (IOException | InterruptedException e) {
+      this.natsConsumer.connect(processor);
+    } catch (SpRuntimeException e) {
       throw new ParseException("Could not connect to Nats broker", e);
     }
 
diff --git 
a/streampipes-extensions/streampipes-sinks-brokers-jvm/src/main/java/org/apache/streampipes/sinks/brokers/jvm/jms/JmsPublisherSink.java
 
b/streampipes-extensions/streampipes-sinks-brokers-jvm/src/main/java/org/apache/streampipes/sinks/brokers/jvm/jms/JmsPublisherSink.java
index a1e0dd649..e41326bd5 100644
--- 
a/streampipes-extensions/streampipes-sinks-brokers-jvm/src/main/java/org/apache/streampipes/sinks/brokers/jvm/jms/JmsPublisherSink.java
+++ 
b/streampipes-extensions/streampipes-sinks-brokers-jvm/src/main/java/org/apache/streampipes/sinks/brokers/jvm/jms/JmsPublisherSink.java
@@ -74,10 +74,12 @@ public class JmsPublisherSink extends StreamPipesDataSink {
     Integer jmsPort = extractor.singleValueParameter(PORT_KEY, Integer.class);
     String topic = extractor.singleValueParameter(TOPIC_KEY, String.class);
 
-    this.publisher = new ActiveMQPublisher();
     JmsTransportProtocol jmsTransportProtocol =
         new JmsTransportProtocol(jmsHost, jmsPort, topic);
-    this.publisher.connect(jmsTransportProtocol);
+
+    this.publisher = new ActiveMQPublisher(jmsTransportProtocol);
+    this.publisher.connect();
+
     if (!this.publisher.isConnected()) {
       throw new SpRuntimeException(
           "Could not connect to JMS server " + jmsHost + " on Port: " + jmsPort
diff --git a/streampipes-integration-tests/pom.xml 
b/streampipes-integration-tests/pom.xml
index f77c4e5b3..cd1959158 100644
--- a/streampipes-integration-tests/pom.xml
+++ b/streampipes-integration-tests/pom.xml
@@ -39,6 +39,12 @@
   </properties>
 
   <dependencies>
+    <dependency>
+      <groupId>org.apache.streampipes</groupId>
+      <artifactId>streampipes-client</artifactId>
+      <version>0.93.0-SNAPSHOT</version>
+      <scope>test</scope>
+    </dependency>
     <dependency>
       <groupId>org.apache.streampipes</groupId>
       <artifactId>streampipes-extensions-management</artifactId>
@@ -55,6 +61,11 @@
       <artifactId>streampipes-messaging-mqtt</artifactId>
       <version>0.93.0-SNAPSHOT</version>
     </dependency>
+    <dependency>
+      <groupId>org.apache.streampipes</groupId>
+      <artifactId>streampipes-messaging-nats</artifactId>
+      <version>0.93.0-SNAPSHOT</version>
+    </dependency>
     <dependency>
       <groupId>org.apache.streampipes</groupId>
       <artifactId>streampipes-connect-adapters-iiot</artifactId>
diff --git 
a/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/adapters/MqttAdapterTester.java
 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/adapters/MqttAdapterTester.java
index 80a54aa32..f7e38f7e1 100644
--- 
a/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/adapters/MqttAdapterTester.java
+++ 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/adapters/MqttAdapterTester.java
@@ -24,6 +24,7 @@ import 
org.apache.streampipes.extensions.api.connect.IAdapterConfiguration;
 import org.apache.streampipes.extensions.api.connect.StreamPipesAdapter;
 import org.apache.streampipes.integration.containers.MosquittoContainer;
 import org.apache.streampipes.integration.containers.MosquittoDevContainer;
+import org.apache.streampipes.integration.utils.Utils;
 import org.apache.streampipes.manager.template.AdapterTemplateHandler;
 import org.apache.streampipes.messaging.mqtt.MqttPublisher;
 import org.apache.streampipes.model.grounding.MqttTransportProtocol;
@@ -36,7 +37,6 @@ import com.fasterxml.jackson.core.JsonProcessingException;
 import com.fasterxml.jackson.databind.ObjectMapper;
 import org.jetbrains.annotations.NotNull;
 
-import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -104,17 +104,7 @@ public class MqttAdapterTester extends AdapterTesterBase {
 
   @Override
   public List<Map<String, Object>> getTestEvents() {
-    List<Map<String, Object>> result = new ArrayList<>();
-
-    for (int i = 0; i < 3; i++) {
-      result.add(
-          Map.of(
-              "timestamp", i,
-              "value", "test-data")
-      );
-    }
-
-    return result;
+    return Utils.getSimpleTestEvents();
   }
 
 
@@ -141,8 +131,8 @@ public class MqttAdapterTester extends AdapterTesterBase {
         mosquittoContainer.getBrokerHost(),
         mosquittoContainer.getBrokerPort(),
         TOPIC);
-    MqttPublisher publisher = new MqttPublisher();
-    publisher.connect(mqttSettings);
+    MqttPublisher publisher = new MqttPublisher(mqttSettings);
+    publisher.connect();
     return publisher;
   }
 
diff --git 
a/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/adapters/PulsarAdapterTester.java
 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/adapters/PulsarAdapterTester.java
index 16bea0d82..177b38a3d 100644
--- 
a/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/adapters/PulsarAdapterTester.java
+++ 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/adapters/PulsarAdapterTester.java
@@ -23,6 +23,7 @@ import 
org.apache.streampipes.extensions.api.connect.IAdapterConfiguration;
 import org.apache.streampipes.extensions.api.connect.StreamPipesAdapter;
 import org.apache.streampipes.integration.containers.PulsarContainer;
 import org.apache.streampipes.integration.containers.PulsarDevContainer;
+import org.apache.streampipes.integration.utils.Utils;
 import org.apache.streampipes.manager.template.AdapterTemplateHandler;
 import org.apache.streampipes.model.staticproperty.StaticPropertyAlternatives;
 import org.apache.streampipes.model.template.PipelineElementTemplate;
@@ -34,7 +35,6 @@ import org.apache.pulsar.client.api.Producer;
 import org.apache.pulsar.client.api.PulsarClient;
 import org.apache.pulsar.client.api.PulsarClientException;
 
-import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
@@ -93,17 +93,7 @@ public class PulsarAdapterTester extends AdapterTesterBase {
 
   @Override
   public List<Map<String, Object>> getTestEvents() {
-    List<Map<String, Object>> result = new ArrayList<>();
-
-    for (int i = 0; i < 3; i++) {
-      result.add(
-          Map.of(
-              "timestamp", i,
-              "value", "test-data")
-      );
-    }
-
-    return result;
+    return Utils.getSimpleTestEvents();
   }
 
   @Override
diff --git 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/client/ClientLiveDataTest.java
similarity index 76%
copy from 
streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
copy to 
streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/client/ClientLiveDataTest.java
index cd24baad5..a456df518 100644
--- 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
+++ 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/client/ClientLiveDataTest.java
@@ -16,13 +16,16 @@
  *
  */
 
-package org.apache.streampipes.messaging;
+package org.apache.streampipes.integration.client;
 
-import org.apache.streampipes.model.grounding.TransportProtocol;
+import org.junit.Test;
 
-public interface SpProtocolDefinition<T extends TransportProtocol> {
+public class ClientLiveDataTest {
 
-  EventConsumer<T> getConsumer();
-
-  EventProducer<T> getProducer();
+  @Test
+  public void testNatsClient() throws Exception {
+    try (var tester = new ClientNatsTester()) {
+      tester.run();
+    }
+  }
 }
diff --git 
a/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/client/ClientLiveDataTesterBase.java
 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/client/ClientLiveDataTesterBase.java
new file mode 100644
index 000000000..0eba57436
--- /dev/null
+++ 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/client/ClientLiveDataTesterBase.java
@@ -0,0 +1,117 @@
+/*
+ * 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.integration.client;
+
+import org.apache.streampipes.client.StreamPipesClient;
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.client.api.live.IConfiguredEventProducer;
+import org.apache.streampipes.client.credentials.StreamPipesApiKeyCredentials;
+import org.apache.streampipes.integration.utils.Utils;
+import org.apache.streampipes.model.SpDataStream;
+import org.apache.streampipes.model.grounding.EventGrounding;
+import org.apache.streampipes.model.grounding.TransportFormat;
+import org.apache.streampipes.model.grounding.TransportProtocol;
+import org.apache.streampipes.vocabulary.MessageFormat;
+
+import org.testcontainers.shaded.com.google.common.collect.Maps;
+
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+
+public abstract class ClientLiveDataTesterBase<T extends TransportProtocol> 
implements AutoCloseable {
+
+  private List<Map<String, Object>> expectedEvents;
+  private int counter;
+
+  public IStreamPipesClient makeStreamPipesClient() {
+    var client = StreamPipesClient.create(
+        "localhost",
+        new StreamPipesApiKeyCredentials("", ""),
+        true
+    );
+    prepareClient(client);
+
+    return client;
+  }
+
+  public void run() throws InterruptedException {
+    startContainer();
+
+    var dataStream = makeDataStream();
+    var client = makeStreamPipesClient();
+
+    var consumer = client.streams().subscribe(dataStream, event -> {
+      assertTrue(Maps.difference(event.getRaw(), 
expectedEvents.get(counter)).areEqual());
+      counter++;
+    });
+
+    var producer = client.streams().getProducer(dataStream);
+
+    expectedEvents = Utils.getSimpleTestEvents();
+
+    publishEvents(producer, expectedEvents);
+
+    // validate that events where send correctly
+    validate(expectedEvents);
+
+    producer.close();
+    consumer.unsubscribe();
+
+  }
+
+  private SpDataStream makeDataStream() {
+    var dataStream = new SpDataStream();
+    dataStream.setEventGrounding(makeEventGrounding());
+
+    return dataStream;
+  }
+
+  private EventGrounding makeEventGrounding() {
+    var grounding = new EventGrounding();
+    grounding.setTransportProtocol(makeProtocol());
+    grounding.setTransportFormats(List.of(new 
TransportFormat(MessageFormat.JSON)));
+
+    return grounding;
+  }
+
+  public void publishEvents(IConfiguredEventProducer producer,
+                            List<Map<String, Object>> expectedEvents) {
+    expectedEvents.forEach(producer::publish);
+  }
+
+  public void validate(List<Map<String, Object>> expectedEvents) throws 
InterruptedException {
+    int retry = 0;
+    while (counter != expectedEvents.size() && retry < 5) {
+      TimeUnit.MILLISECONDS.sleep(1000);
+      retry++;
+    }
+
+    assertEquals(expectedEvents.size(), counter);
+  }
+
+  public abstract void startContainer();
+
+  public abstract T makeProtocol();
+
+  public abstract void prepareClient(IStreamPipesClient client);
+}
diff --git 
a/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/client/ClientNatsTester.java
 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/client/ClientNatsTester.java
new file mode 100644
index 000000000..d2fa57786
--- /dev/null
+++ 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/client/ClientNatsTester.java
@@ -0,0 +1,64 @@
+/*
+ * 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.integration.client;
+
+import org.apache.streampipes.client.api.IStreamPipesClient;
+import org.apache.streampipes.integration.containers.NatsContainer;
+import org.apache.streampipes.integration.containers.NatsDevContainer;
+import org.apache.streampipes.messaging.nats.SpNatsProtocolFactory;
+import org.apache.streampipes.model.grounding.NatsTransportProtocol;
+import org.apache.streampipes.model.grounding.SimpleTopicDefinition;
+
+import java.util.Objects;
+
+public class ClientNatsTester extends 
ClientLiveDataTesterBase<NatsTransportProtocol> {
+
+  private NatsContainer natsContainer;
+
+  @Override
+  public void startContainer() {
+    if (Objects.equals(System.getenv("TEST_MODE"), "dev")) {
+      natsContainer = new NatsDevContainer();
+    } else {
+      natsContainer = new NatsContainer();
+    }
+
+    natsContainer.start();
+  }
+
+  @Override
+  public NatsTransportProtocol makeProtocol() {
+    var protocol = new NatsTransportProtocol();
+    protocol.setBrokerHostname(natsContainer.getBrokerHost());
+    protocol.setPort(natsContainer.getBrokerPort());
+    protocol.setTopicDefinition(new SimpleTopicDefinition("test-topic"));
+
+    return protocol;
+  }
+
+  @Override
+  public void prepareClient(IStreamPipesClient client) {
+    client.registerProtocol(new SpNatsProtocolFactory());
+  }
+
+  @Override
+  public void close() throws Exception {
+    natsContainer.close();
+  }
+}
diff --git 
a/streampipes-client/src/main/java/org/apache/streampipes/client/live/KafkaConfig.java
 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/containers/NatsContainer.java
similarity index 52%
rename from 
streampipes-client/src/main/java/org/apache/streampipes/client/live/KafkaConfig.java
rename to 
streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/containers/NatsContainer.java
index ca25c1328..76a0f833f 100644
--- 
a/streampipes-client/src/main/java/org/apache/streampipes/client/live/KafkaConfig.java
+++ 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/containers/NatsContainer.java
@@ -15,31 +15,35 @@
  * limitations under the License.
  *
  */
-package org.apache.streampipes.client.live;
 
-import org.apache.streampipes.client.api.live.IKafkaConfig;
+package org.apache.streampipes.integration.containers;
 
-public class KafkaConfig implements IKafkaConfig {
+import org.testcontainers.containers.GenericContainer;
+import org.testcontainers.containers.wait.strategy.LogMessageWaitStrategy;
 
-  private String kafkaHost;
-  private Integer kafkaPort;
+public class NatsContainer extends GenericContainer<NatsContainer> {
 
-  private KafkaConfig(String kafkaHost, Integer kafkaPort) {
-    this.kafkaHost = kafkaHost;
-    this.kafkaPort = kafkaPort;
+  protected static final int NATS_PORT = 4222;
+
+  public NatsContainer() {
+    super("nats:latest");
+  }
+
+  public void start() {
+    this.withExposedPorts(NATS_PORT);
+    this.waitingFor(new LogMessageWaitStrategy().withRegEx(".*Server is 
ready.*"));
+    super.start();
   }
 
-  public static KafkaConfig create(String kafkaHost, Integer kafkaPort) {
-    return new KafkaConfig(kafkaHost, kafkaPort);
+  public String getBrokerHost() {
+    return getHost();
   }
 
-  @Override
-  public String getKafkaHost() {
-    return kafkaHost;
+  public Integer getBrokerPort() {
+    return getMappedPort(NATS_PORT);
   }
 
-  @Override
-  public Integer getKafkaPort() {
-    return kafkaPort;
+  public String getBrokerUrl() {
+    return getBrokerHost() + ":" + getBrokerPort();
   }
 }
diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IKafkaConfig.java
 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/containers/NatsDevContainer.java
similarity index 81%
rename from 
streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IKafkaConfig.java
rename to 
streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/containers/NatsDevContainer.java
index 6e70c06fe..36e966303 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/live/IKafkaConfig.java
+++ 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/containers/NatsDevContainer.java
@@ -16,10 +16,12 @@
  *
  */
 
-package org.apache.streampipes.client.api.live;
+package org.apache.streampipes.integration.containers;
 
-public interface IKafkaConfig {
-  String getKafkaHost();
+public class NatsDevContainer extends NatsContainer {
 
-  Integer getKafkaPort();
+  @Override
+  public String getBrokerHost() {
+    return "localhost";
+  }
 }
diff --git 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/utils/Utils.java
similarity index 66%
copy from 
streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
copy to 
streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/utils/Utils.java
index b1576ce67..b27fce623 100644
--- 
a/streampipes-client-api/src/main/java/org/apache/streampipes/client/api/config/IStreamPipesClientConfig.java
+++ 
b/streampipes-integration-tests/src/test/java/org/apache/streampipes/integration/utils/Utils.java
@@ -16,20 +16,25 @@
  *
  */
 
-package org.apache.streampipes.client.api.config;
-
-import org.apache.streampipes.dataformat.SpDataFormatFactory;
-
-import com.fasterxml.jackson.databind.ObjectMapper;
+package org.apache.streampipes.integration.utils;
 
+import java.util.ArrayList;
 import java.util.List;
+import java.util.Map;
 
-public interface IStreamPipesClientConfig {
-  ObjectMapper getSerializer();
+public class Utils {
 
-  void addDataFormat(SpDataFormatFactory spDataFormatFactory);
+  public static List<Map<String, Object>> getSimpleTestEvents() {
+    List<Map<String, Object>> result = new ArrayList<>();
 
-  List<SpDataFormatFactory> getRegisteredDataFormats();
+    for (int i = 0; i < 3; i++) {
+      result.add(
+          Map.of(
+              "timestamp", i,
+              "value", "test-data")
+      );
+    }
 
-  ClientConnectionUrlResolver getConnectionConfig();
+    return result;
+  }
 }
diff --git 
a/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQConnectionProvider.java
 
b/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQConnectionProvider.java
index 4de17d167..5404c591a 100644
--- 
a/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQConnectionProvider.java
+++ 
b/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQConnectionProvider.java
@@ -18,6 +18,8 @@
 
 package org.apache.streampipes.messaging.jms;
 
+import org.apache.streampipes.model.grounding.JmsTransportProtocol;
+
 import org.apache.activemq.ActiveMQConnectionFactory;
 
 import javax.jms.Connection;
@@ -25,6 +27,12 @@ import javax.jms.JMSException;
 
 public abstract class ActiveMQConnectionProvider {
 
+  protected JmsTransportProtocol protocol;
+
+  public ActiveMQConnectionProvider(JmsTransportProtocol protocol) {
+    this.protocol = protocol;
+  }
+
   protected Connection startJmsConnection(String url) {
     try {
       ActiveMQConnectionFactory connectionFactory = new 
ActiveMQConnectionFactory(url);
diff --git 
a/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQConsumer.java
 
b/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQConsumer.java
index f043f785b..3cef5b98f 100644
--- 
a/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQConsumer.java
+++ 
b/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQConsumer.java
@@ -34,7 +34,7 @@ import javax.jms.Session;
 import java.io.Serializable;
 
 public class ActiveMQConsumer extends ActiveMQConnectionProvider implements
-    EventConsumer<JmsTransportProtocol>,
+    EventConsumer,
     AutoCloseable, Serializable {
 
   private Session session;
@@ -43,6 +43,10 @@ public class ActiveMQConsumer extends 
ActiveMQConnectionProvider implements
 
   private Boolean connected = false;
 
+  public ActiveMQConsumer(JmsTransportProtocol protocol) {
+    super(protocol);
+  }
+
   private void initListener() {
     try {
       consumer.setMessageListener(message -> {
@@ -58,15 +62,15 @@ public class ActiveMQConsumer extends 
ActiveMQConnectionProvider implements
   }
 
   @Override
-  public void connect(JmsTransportProtocol protocolSettings, 
InternalEventProcessor<byte[]>
+  public void connect(InternalEventProcessor<byte[]>
       eventProcessor) throws SpRuntimeException {
-    String url = ActiveMQUtils.makeActiveMqUrl(protocolSettings);
+    String url = ActiveMQUtils.makeActiveMqUrl(protocol);
 
     try {
       this.eventProcessor = eventProcessor;
       session = startJmsConnection(url).createSession(false, 
Session.AUTO_ACKNOWLEDGE);
       consumer = session.createConsumer(session.createTopic(
-          protocolSettings.getTopicDefinition().getActualTopicName())
+          protocol.getTopicDefinition().getActualTopicName())
       );
       initListener();
       this.connected = true;
diff --git 
a/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQPublisher.java
 
b/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQPublisher.java
index 374f504b2..f5443a95e 100644
--- 
a/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQPublisher.java
+++ 
b/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/ActiveMQPublisher.java
@@ -21,7 +21,6 @@ package org.apache.streampipes.messaging.jms;
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
 import org.apache.streampipes.messaging.EventProducer;
 import org.apache.streampipes.model.grounding.JmsTransportProtocol;
-import org.apache.streampipes.model.grounding.SimpleTopicDefinition;
 
 import org.apache.activemq.ActiveMQConnectionFactory;
 import org.slf4j.Logger;
@@ -36,7 +35,7 @@ import javax.jms.MessageProducer;
 import javax.jms.Session;
 
 
-public class ActiveMQPublisher implements EventProducer<JmsTransportProtocol> {
+public class ActiveMQPublisher extends ActiveMQConnectionProvider implements 
EventProducer {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(ActiveMQPublisher.class);
 
@@ -46,43 +45,14 @@ public class ActiveMQPublisher implements 
EventProducer<JmsTransportProtocol> {
 
   private boolean connected = false;
 
-  public ActiveMQPublisher() {
-
-  }
-
-  @Deprecated
-  public ActiveMQPublisher(String url, String topic) {
-    JmsTransportProtocol protocol = new JmsTransportProtocol();
-    protocol.setBrokerHostname(url.substring(0, url.lastIndexOf(":")));
-    protocol.setPort(Integer.parseInt(url.substring(url.lastIndexOf(":") + 1, 
url.length())));
-    protocol.setTopicDefinition(new SimpleTopicDefinition(topic));
-    try {
-      connect(protocol);
-    } catch (SpRuntimeException e) {
-      e.printStackTrace();
-    }
-  }
-
-  public ActiveMQPublisher(String host, int port, String topic) {
-    JmsTransportProtocol protocol = new JmsTransportProtocol();
-    protocol.setBrokerHostname(host);
-    protocol.setPort(port);
-    protocol.setTopicDefinition(new SimpleTopicDefinition(topic));
-    try {
-      connect(protocol);
-    } catch (SpRuntimeException e) {
-      e.printStackTrace();
-    }
-  }
-
-  public void sendText(String message) throws JMSException {
-    publish(message.getBytes());
+  public ActiveMQPublisher(JmsTransportProtocol protocol) {
+    super(protocol);
   }
 
   @Override
-  public void connect(JmsTransportProtocol protocolSettings) throws 
SpRuntimeException {
+  public void connect() throws SpRuntimeException {
 
-    String url = ActiveMQUtils.makeActiveMqUrl(protocolSettings);
+    String url = ActiveMQUtils.makeActiveMqUrl(protocol);
     ConnectionFactory connectionFactory = new ActiveMQConnectionFactory(url);
 
     boolean co = false;
@@ -98,7 +68,7 @@ public class ActiveMQPublisher implements 
EventProducer<JmsTransportProtocol> {
     try {
       this.session = connection
           .createSession(false, Session.AUTO_ACKNOWLEDGE);
-      this.producer = 
session.createProducer(session.createTopic(protocolSettings
+      this.producer = session.createProducer(session.createTopic(protocol
           .getTopicDefinition()
           .getActualTopicName()));
       this.producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT);
@@ -106,9 +76,8 @@ public class ActiveMQPublisher implements 
EventProducer<JmsTransportProtocol> {
       this.connected = true;
     } catch (JMSException e) {
       throw new SpRuntimeException("could not connect to activemq broker. 
Broker: '"
-          + protocolSettings.getBrokerHostname() + "' Port: " + 
protocolSettings.getPort());
+          + protocol.getBrokerHostname() + "' Port: " + protocol.getPort());
     }
-
   }
 
   @Override
@@ -130,9 +99,7 @@ public class ActiveMQPublisher implements 
EventProducer<JmsTransportProtocol> {
       session.close();
       connection.close();
       this.connected = false;
-      //logger.info("ActiveMQ connection closed successfully.");
     } catch (JMSException e) {
-      //logger.warn("Could not close ActiveMQ connection.");
       throw new SpRuntimeException("could not disconnect from activemq 
broker");
     }
   }
diff --git 
a/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/SpJmsProtocol.java
 
b/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/SpJmsProtocol.java
index dcdf03c6f..b8f7f5aba 100644
--- 
a/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/SpJmsProtocol.java
+++ 
b/streampipes-messaging-jms/src/main/java/org/apache/streampipes/messaging/jms/SpJmsProtocol.java
@@ -24,21 +24,14 @@ import 
org.apache.streampipes.model.grounding.JmsTransportProtocol;
 
 public class SpJmsProtocol implements 
SpProtocolDefinition<JmsTransportProtocol> {
 
-  private final EventConsumer<JmsTransportProtocol> jmsConsumer;
-  private final EventProducer<JmsTransportProtocol> jmsProducer;
-
-  public SpJmsProtocol() {
-    this.jmsConsumer = new ActiveMQConsumer();
-    this.jmsProducer = new ActiveMQPublisher();
-  }
 
   @Override
-  public EventConsumer<JmsTransportProtocol> getConsumer() {
-    return jmsConsumer;
+  public EventConsumer getConsumer(JmsTransportProtocol transportProtocol) {
+    return new ActiveMQConsumer(transportProtocol);
   }
 
   @Override
-  public EventProducer<JmsTransportProtocol> getProducer() {
-    return jmsProducer;
+  public EventProducer getProducer(JmsTransportProtocol transportProtocol) {
+    return new ActiveMQPublisher(transportProtocol);
   }
 }
diff --git 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaConsumer.java
 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaConsumer.java
index 1411b9c03..e999cfbc4 100644
--- 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaConsumer.java
+++ 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaConsumer.java
@@ -43,12 +43,12 @@ import java.util.List;
 import java.util.Properties;
 import java.util.regex.Pattern;
 
-public class SpKafkaConsumer implements EventConsumer<KafkaTransportProtocol>, 
Runnable,
+public class SpKafkaConsumer implements EventConsumer, Runnable,
     Serializable {
 
   private String topic;
   private InternalEventProcessor<byte[]> eventProcessor;
-  private KafkaTransportProtocol protocol;
+  private final KafkaTransportProtocol protocol;
   private volatile boolean isRunning;
   private Boolean patternTopic = false;
 
@@ -56,8 +56,8 @@ public class SpKafkaConsumer implements 
EventConsumer<KafkaTransportProtocol>, R
 
   private static final Logger LOG = 
LoggerFactory.getLogger(SpKafkaConsumer.class);
 
-  public SpKafkaConsumer() {
-
+  public SpKafkaConsumer(KafkaTransportProtocol protocol) {
+    this.protocol = protocol;
   }
 
   public SpKafkaConsumer(KafkaTransportProtocol protocol,
@@ -119,15 +119,13 @@ public class SpKafkaConsumer implements 
EventConsumer<KafkaTransportProtocol>, R
   }
 
   @Override
-  public void connect(KafkaTransportProtocol protocol, 
InternalEventProcessor<byte[]>
-      eventProcessor)
-      throws SpRuntimeException {
+  public void connect(InternalEventProcessor<byte[]> eventProcessor) throws 
SpRuntimeException {
     LOG.info("Kafka consumer: Connecting to " + 
protocol.getTopicDefinition().getActualTopicName());
     if (protocol.getTopicDefinition() instanceof WildcardTopicDefinition) {
       this.patternTopic = true;
     }
     this.eventProcessor = eventProcessor;
-    this.protocol = protocol;
+
     this.topic = protocol.getTopicDefinition().getActualTopicName();
     this.isRunning = true;
 
diff --git 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProducer.java
 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProducer.java
index 9f0d262ee..7938ecf65 100644
--- 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProducer.java
+++ 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProducer.java
@@ -45,7 +45,7 @@ import java.util.Map;
 import java.util.Properties;
 import java.util.concurrent.ExecutionException;
 
-public class SpKafkaProducer implements EventProducer<KafkaTransportProtocol>, 
Serializable {
+public class SpKafkaProducer implements EventProducer, Serializable {
 
 
   private static final String COLON = ":";
@@ -53,12 +53,14 @@ public class SpKafkaProducer implements 
EventProducer<KafkaTransportProtocol>, S
   private String brokerUrl;
   private String topic;
   private Producer<String, byte[]> producer;
+  private KafkaTransportProtocol protocol;
 
   private boolean connected = false;
 
   private static final Logger LOG = 
LoggerFactory.getLogger(SpKafkaProducer.class);
 
-  public SpKafkaProducer() {
+  public SpKafkaProducer(KafkaTransportProtocol protocol) {
+    this.protocol = protocol;
   }
 
   // TODO backwards compatibility, remove later
@@ -90,7 +92,7 @@ public class SpKafkaProducer implements 
EventProducer<KafkaTransportProtocol>, S
   }
 
   @Override
-  public void connect(KafkaTransportProtocol protocol) {
+  public void connect() {
     LOG.info("Kafka producer: Connecting to " + 
protocol.getTopicDefinition().getActualTopicName());
     this.brokerUrl = protocol.getBrokerHostname() + ":" + 
protocol.getKafkaPort();
     this.topic = protocol.getTopicDefinition().getActualTopicName();
diff --git 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProtocol.java
 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProtocol.java
index 0a921092d..a75ebb775 100644
--- 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProtocol.java
+++ 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProtocol.java
@@ -25,21 +25,13 @@ import 
org.apache.streampipes.model.grounding.KafkaTransportProtocol;
 
 public class SpKafkaProtocol implements 
SpProtocolDefinition<KafkaTransportProtocol> {
 
-  private final EventConsumer<KafkaTransportProtocol> kafkaConsumer;
-  private final EventProducer<KafkaTransportProtocol> kafkaProducer;
-
-  public SpKafkaProtocol() {
-    this.kafkaConsumer = new SpKafkaConsumer();
-    this.kafkaProducer = new SpKafkaProducer();
-  }
-
   @Override
-  public EventConsumer<KafkaTransportProtocol> getConsumer() {
-    return kafkaConsumer;
+  public EventConsumer getConsumer(KafkaTransportProtocol transportProtocol) {
+    return new SpKafkaConsumer(transportProtocol);
   }
 
   @Override
-  public EventProducer<KafkaTransportProtocol> getProducer() {
-    return kafkaProducer;
+  public EventProducer getProducer(KafkaTransportProtocol transportProtocol) {
+    return new SpKafkaProducer(transportProtocol);
   }
 }
diff --git 
a/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/AbstractMqttConnector.java
 
b/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/AbstractMqttConnector.java
index c5cad2b6d..b5270b508 100644
--- 
a/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/AbstractMqttConnector.java
+++ 
b/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/AbstractMqttConnector.java
@@ -28,6 +28,12 @@ public class AbstractMqttConnector {
   protected BlockingConnection connection;
   protected boolean connected = false;
 
+  protected final MqttTransportProtocol protocol;
+
+  public AbstractMqttConnector(MqttTransportProtocol protocol) {
+    this.protocol = protocol;
+  }
+
   protected void createBrokerConnection(MqttTransportProtocol 
protocolSettings) throws Exception {
     this.mqtt = new MQTT();
     this.mqtt.setHost(makeBrokerUrl(protocolSettings));
diff --git 
a/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/MqttConsumer.java
 
b/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/MqttConsumer.java
index d19d4b575..dbe3b3e2f 100644
--- 
a/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/MqttConsumer.java
+++ 
b/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/MqttConsumer.java
@@ -29,16 +29,20 @@ import org.fusesource.mqtt.client.Topic;
 import java.io.Serializable;
 
 public class MqttConsumer extends AbstractMqttConnector implements
-    EventConsumer<MqttTransportProtocol>,
+    EventConsumer,
     AutoCloseable, Serializable {
 
+  public MqttConsumer(MqttTransportProtocol protocol) {
+    super(protocol);
+  }
+
   @Override
-  public void connect(MqttTransportProtocol protocolSettings, 
InternalEventProcessor<byte[]> eventProcessor)
+  public void connect(InternalEventProcessor<byte[]> eventProcessor)
       throws SpRuntimeException {
 
     try {
-      this.createBrokerConnection(protocolSettings);
-      Topic[] topics = {new 
Topic(protocolSettings.getTopicDefinition().getActualTopicName(), 
QoS.AT_LEAST_ONCE)};
+      this.createBrokerConnection(protocol);
+      Topic[] topics = {new 
Topic(protocol.getTopicDefinition().getActualTopicName(), QoS.AT_LEAST_ONCE)};
       connection.subscribe(topics);
       new Thread(new ConsumerThread(eventProcessor)).start();
 
diff --git 
a/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/MqttPublisher.java
 
b/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/MqttPublisher.java
index cf4f8356e..d9af23ec5 100644
--- 
a/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/MqttPublisher.java
+++ 
b/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/MqttPublisher.java
@@ -25,17 +25,21 @@ import org.fusesource.mqtt.client.QoS;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-public class MqttPublisher extends AbstractMqttConnector implements 
EventProducer<MqttTransportProtocol> {
+public class MqttPublisher extends AbstractMqttConnector implements 
EventProducer {
 
   private static final Logger LOG = 
LoggerFactory.getLogger(MqttPublisher.class);
 
   private String currentTopic;
 
+  public MqttPublisher(MqttTransportProtocol protocol) {
+    super(protocol);
+  }
+
   @Override
-  public void connect(MqttTransportProtocol protocolSettings) throws 
SpRuntimeException {
+  public void connect() throws SpRuntimeException {
     try {
-      this.createBrokerConnection(protocolSettings);
-      this.currentTopic = 
protocolSettings.getTopicDefinition().getActualTopicName();
+      this.createBrokerConnection(protocol);
+      this.currentTopic = protocol.getTopicDefinition().getActualTopicName();
     } catch (Exception e) {
       throw new SpRuntimeException(e);
     }
diff --git 
a/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/SpMqttProtocol.java
 
b/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/SpMqttProtocol.java
index ccda09075..a80caf378 100644
--- 
a/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/SpMqttProtocol.java
+++ 
b/streampipes-messaging-mqtt/src/main/java/org/apache/streampipes/messaging/mqtt/SpMqttProtocol.java
@@ -24,21 +24,13 @@ import 
org.apache.streampipes.model.grounding.MqttTransportProtocol;
 
 public class SpMqttProtocol implements 
SpProtocolDefinition<MqttTransportProtocol> {
 
-  private final EventConsumer<MqttTransportProtocol> mqttConsumer;
-  private final EventProducer<MqttTransportProtocol> mqttProducer;
-
-  public SpMqttProtocol() {
-    this.mqttConsumer = new MqttConsumer();
-    this.mqttProducer = new MqttPublisher();
-  }
-
   @Override
-  public EventConsumer<MqttTransportProtocol> getConsumer() {
-    return this.mqttConsumer;
+  public EventConsumer getConsumer(MqttTransportProtocol transportProtocol) {
+    return new MqttConsumer(transportProtocol);
   }
 
   @Override
-  public EventProducer<MqttTransportProtocol> getProducer() {
-    return this.mqttProducer;
+  public EventProducer getProducer(MqttTransportProtocol transportProtocol) {
+    return new MqttPublisher(transportProtocol);
   }
 }
diff --git 
a/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/AbstractNatsConnector.java
 
b/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/AbstractNatsConnector.java
index 5253cdf24..9b7d1f76d 100644
--- 
a/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/AbstractNatsConnector.java
+++ 
b/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/AbstractNatsConnector.java
@@ -45,6 +45,14 @@ public abstract class AbstractNatsConnector {
     this.subject = protocol.getTopicDefinition().getActualTopicName();
   }
 
+  protected NatsConfig makeNatsConfig(NatsTransportProtocol protocol) {
+    var natsConfig = new NatsConfig();
+    natsConfig.setNatsUrls(makeBrokerUrl(protocol));
+    natsConfig.setSubject(protocol.getTopicDefinition().getActualTopicName());
+
+    return natsConfig;
+  }
+
   protected void disconnect() throws InterruptedException, TimeoutException {
     natsConnection.flush(Duration.ofMillis(50));
     natsConnection.close();
diff --git 
a/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/NatsConsumer.java
 
b/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/NatsConsumer.java
index ee02fab77..32eaeb2b9 100644
--- 
a/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/NatsConsumer.java
+++ 
b/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/NatsConsumer.java
@@ -31,25 +31,33 @@ import io.nats.client.Subscription;
 import java.io.IOException;
 import java.util.concurrent.TimeoutException;
 
-public class NatsConsumer extends AbstractNatsConnector implements 
EventConsumer<NatsTransportProtocol> {
+public class NatsConsumer extends AbstractNatsConnector implements 
EventConsumer {
 
   private Dispatcher dispatcher;
   private Subscription subscription;
+  private NatsConfig natsConfig;
+
+  public NatsConsumer(NatsTransportProtocol protocol) {
+    this.natsConfig = makeNatsConfig(protocol);
+  }
+
+  public NatsConsumer(NatsConfig natsConfig) {
+    this.natsConfig = natsConfig;
+  }
 
   public void connect(NatsConfig natsConfig,
                       InternalEventProcessor<byte[]> eventProcessor) throws 
IOException, InterruptedException {
-    makeBrokerConnection(natsConfig);
-    createSubscription(eventProcessor);
+    this.natsConfig = natsConfig;
+    connect(eventProcessor);
   }
 
   @Override
-  public void connect(NatsTransportProtocol protocolSettings,
-                      InternalEventProcessor<byte[]> eventProcessor) throws 
SpRuntimeException {
+  public void connect(InternalEventProcessor<byte[]> eventProcessor) throws 
SpRuntimeException {
     try {
-      makeBrokerConnection(protocolSettings);
+      makeBrokerConnection(natsConfig);
       createSubscription(eventProcessor);
     } catch (IOException | InterruptedException e) {
-      e.printStackTrace();
+      throw new SpRuntimeException(e);
     }
   }
 
diff --git 
a/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/NatsPublisher.java
 
b/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/NatsPublisher.java
index 06e7cc732..7d0551fae 100644
--- 
a/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/NatsPublisher.java
+++ 
b/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/NatsPublisher.java
@@ -27,12 +27,18 @@ import io.nats.client.Connection;
 import java.io.IOException;
 import java.util.concurrent.TimeoutException;
 
-public class NatsPublisher extends AbstractNatsConnector implements 
EventProducer<NatsTransportProtocol> {
+public class NatsPublisher extends AbstractNatsConnector implements 
EventProducer {
+
+  private final NatsTransportProtocol protocol;
+
+  public NatsPublisher(NatsTransportProtocol protocol) {
+    this.protocol = protocol;
+  }
 
   @Override
-  public void connect(NatsTransportProtocol protocolSettings) throws 
SpRuntimeException {
+  public void connect() throws SpRuntimeException {
     try {
-      makeBrokerConnection(protocolSettings);
+      makeBrokerConnection(protocol);
     } catch (IOException | InterruptedException e) {
       e.printStackTrace();
     }
diff --git 
a/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/SpNatsProtocol.java
 
b/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/SpNatsProtocol.java
index bf3abb3de..052c33749 100644
--- 
a/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/SpNatsProtocol.java
+++ 
b/streampipes-messaging-nats/src/main/java/org/apache/streampipes/messaging/nats/SpNatsProtocol.java
@@ -25,21 +25,13 @@ import 
org.apache.streampipes.model.grounding.NatsTransportProtocol;
 
 public class SpNatsProtocol implements 
SpProtocolDefinition<NatsTransportProtocol> {
 
-  private final EventConsumer<NatsTransportProtocol> natsConsumer;
-  private final EventProducer<NatsTransportProtocol> natsProducer;
-
-  public SpNatsProtocol() {
-    this.natsConsumer = new NatsConsumer();
-    this.natsProducer = new NatsPublisher();
-  }
-
   @Override
-  public EventConsumer<NatsTransportProtocol> getConsumer() {
-    return this.natsConsumer;
+  public EventConsumer getConsumer(NatsTransportProtocol transportProtocol) {
+    return new NatsConsumer(transportProtocol);
   }
 
   @Override
-  public EventProducer<NatsTransportProtocol> getProducer() {
-    return this.natsProducer;
+  public EventProducer getProducer(NatsTransportProtocol transportProtocol) {
+    return new NatsPublisher(transportProtocol);
   }
 }
diff --git 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/EventConsumer.java
 
b/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/EventConsumer.java
index dec6f0258..c95e1a789 100644
--- 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/EventConsumer.java
+++ 
b/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/EventConsumer.java
@@ -19,12 +19,10 @@
 package org.apache.streampipes.messaging;
 
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
-import org.apache.streampipes.model.grounding.TransportProtocol;
 
-public interface EventConsumer<T extends TransportProtocol> {
+public interface EventConsumer {
 
-  void connect(T protocolSettings, InternalEventProcessor<byte[]> 
eventProcessor) throws
-      SpRuntimeException;
+  void connect(InternalEventProcessor<byte[]> eventProcessor) throws 
SpRuntimeException;
 
   void disconnect() throws SpRuntimeException;
 
diff --git 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/EventProducer.java
 
b/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/EventProducer.java
index 20993d932..68f9a540a 100644
--- 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/EventProducer.java
+++ 
b/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/EventProducer.java
@@ -19,13 +19,12 @@
 package org.apache.streampipes.messaging;
 
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
-import org.apache.streampipes.model.grounding.TransportProtocol;
 
 import java.io.Serializable;
 
-public interface EventProducer<T extends TransportProtocol> extends 
Serializable {
+public interface EventProducer extends Serializable {
 
-  void connect(T protocolSettings) throws SpRuntimeException;
+  void connect() throws SpRuntimeException;
 
   void publish(byte[] event);
 
diff --git 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
 
b/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
index cd24baad5..3f4e1530a 100644
--- 
a/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
+++ 
b/streampipes-messaging/src/main/java/org/apache/streampipes/messaging/SpProtocolDefinition.java
@@ -22,7 +22,8 @@ import 
org.apache.streampipes.model.grounding.TransportProtocol;
 
 public interface SpProtocolDefinition<T extends TransportProtocol> {
 
-  EventConsumer<T> getConsumer();
+  EventConsumer getConsumer(T transportProtocol);
+
+  EventProducer getProducer(T transportProtocol);
 
-  EventProducer<T> getProducer();
 }
diff --git 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/runtime/PipelineElementRuntimeInfoFetcher.java
 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/runtime/PipelineElementRuntimeInfoFetcher.java
index 2762bad07..9413582db 100644
--- 
a/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/runtime/PipelineElementRuntimeInfoFetcher.java
+++ 
b/streampipes-pipeline-management/src/main/java/org/apache/streampipes/manager/runtime/PipelineElementRuntimeInfoFetcher.java
@@ -20,6 +20,7 @@ package org.apache.streampipes.manager.runtime;
 import org.apache.streampipes.commons.environment.Environment;
 import org.apache.streampipes.commons.environment.Environments;
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
+import org.apache.streampipes.messaging.EventConsumer;
 import org.apache.streampipes.messaging.jms.ActiveMQConsumer;
 import org.apache.streampipes.messaging.kafka.SpKafkaConsumer;
 import org.apache.streampipes.messaging.mqtt.MqttConsumer;
@@ -31,20 +32,16 @@ import 
org.apache.streampipes.model.grounding.MqttTransportProtocol;
 import org.apache.streampipes.model.grounding.NatsTransportProtocol;
 import org.apache.streampipes.model.grounding.TransportFormat;
 
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
 
 public enum PipelineElementRuntimeInfoFetcher {
   INSTANCE;
 
-  Logger logger = 
LoggerFactory.getLogger(PipelineElementRuntimeInfoFetcher.class);
-
   private static final int FETCH_INTERVAL_MS = 300;
   private final Map<String, SpDataFormatConverter> converterMap;
-  private Environment env;
+  private final Environment env;
 
   PipelineElementRuntimeInfoFetcher() {
     this.converterMap = new HashMap<>();
@@ -91,7 +88,7 @@ public enum PipelineElementRuntimeInfoFetcher {
     long timeout = 0;
     while (result[0] == null && timeout < 6000) {
       try {
-        Thread.sleep(FETCH_INTERVAL_MS);
+        TimeUnit.MILLISECONDS.sleep(FETCH_INTERVAL_MS);
         timeout = timeout + 300;
       } catch (InterruptedException e) {
         e.printStackTrace();
@@ -101,39 +98,25 @@ public enum PipelineElementRuntimeInfoFetcher {
 
   private String getLatestEventFromJms(JmsTransportProtocol protocol,
                                        SpDataFormatConverter converter) throws 
SpRuntimeException {
-    final String[] result = {null};
-    ActiveMQConsumer consumer = new ActiveMQConsumer();
-    consumer.connect(protocol, event -> {
-      result[0] = converter.convert(event);
-      consumer.disconnect();
-    });
-
-    waitForEvent(result);
-
-    return result[0];
+    return getLatestEvent(new ActiveMQConsumer(protocol), converter);
   }
 
   private String getLatestEventFromMqtt(MqttTransportProtocol protocol,
                                         SpDataFormatConverter converter) 
throws SpRuntimeException {
-    final String[] result = {null};
-    MqttConsumer mqttConsumer = new MqttConsumer();
-    mqttConsumer.connect(protocol, event -> {
-      result[0] = converter.convert(event);
-      mqttConsumer.disconnect();
-    });
-
-    waitForEvent(result);
-
-    return result[0];
+    return getLatestEvent(new MqttConsumer(protocol), converter);
   }
 
   private String getLatestEventFromNats(NatsTransportProtocol protocol,
                                         SpDataFormatConverter converter) 
throws SpRuntimeException {
+    return getLatestEvent(new NatsConsumer(protocol), converter);
+  }
+
+  private String getLatestEvent(EventConsumer consumer,
+                                SpDataFormatConverter converter) {
     final String[] result = {null};
-    NatsConsumer natsConsumer = new NatsConsumer();
-    natsConsumer.connect(protocol, event -> {
+    consumer.connect(event -> {
       result[0] = converter.convert(event);
-      natsConsumer.disconnect();
+      consumer.disconnect();
     });
 
     waitForEvent(result);
diff --git 
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpInputCollector.java
 
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpInputCollector.java
index 4d7051077..890f16da2 100644
--- 
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpInputCollector.java
+++ 
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpInputCollector.java
@@ -21,6 +21,7 @@ package org.apache.streampipes.wrapper.standalone.routing;
 import org.apache.streampipes.commons.exceptions.SpRuntimeException;
 import org.apache.streampipes.extensions.api.pe.routing.RawDataProcessor;
 import org.apache.streampipes.extensions.api.pe.routing.SpInputCollector;
+import org.apache.streampipes.messaging.EventConsumer;
 import org.apache.streampipes.messaging.InternalEventProcessor;
 import org.apache.streampipes.model.grounding.TransportFormat;
 import org.apache.streampipes.model.grounding.TransportProtocol;
@@ -32,10 +33,13 @@ public class StandaloneSpInputCollector<T extends 
TransportProtocol> extends
     InternalEventProcessor<byte[]>, SpInputCollector {
 
   private final Boolean singletonEngine;
+  private final EventConsumer consumer;
 
-  public StandaloneSpInputCollector(T protocol, TransportFormat format,
+  public StandaloneSpInputCollector(T protocol,
+                                    TransportFormat format,
                                     Boolean singletonEngine) throws 
SpRuntimeException {
     super(protocol, format);
+    this.consumer = protocolDefinition.getConsumer(protocol);
     this.singletonEngine = singletonEngine;
   }
 
@@ -54,16 +58,16 @@ public class StandaloneSpInputCollector<T extends 
TransportProtocol> extends
 
   @Override
   public void connect() throws SpRuntimeException {
-    if (!protocolDefinition.getConsumer().isConnected()) {
-      protocolDefinition.getConsumer().connect(transportProtocol, this);
+    if (!consumer.isConnected()) {
+      consumer.connect(this);
     }
   }
 
   @Override
   public void disconnect() throws SpRuntimeException {
-    if (protocolDefinition.getConsumer().isConnected()) {
+    if (consumer.isConnected()) {
       if (consumers.size() == 0) {
-        protocolDefinition.getConsumer().disconnect();
+        consumer.disconnect();
         ProtocolManager.removeInputCollector(transportProtocol);
       }
     }
diff --git 
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpOutputCollector.java
 
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpOutputCollector.java
index 7cca11e35..c76df7cc6 100644
--- 
a/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpOutputCollector.java
+++ 
b/streampipes-wrapper-standalone/src/main/java/org/apache/streampipes/wrapper/standalone/routing/StandaloneSpOutputCollector.java
@@ -42,14 +42,14 @@ public class StandaloneSpOutputCollector<T extends 
TransportProtocol> extends
 
   private static final Logger LOG = 
LoggerFactory.getLogger(StandaloneSpOutputCollector.class);
 
-  private final EventProducer<T> producer;
+  private final EventProducer producer;
   private final String resourceId;
 
   public StandaloneSpOutputCollector(T protocol,
                                      TransportFormat format,
                                      String resourceId) throws 
SpRuntimeException {
     super(protocol, format);
-    this.producer = protocolDefinition.getProducer();
+    this.producer = protocolDefinition.getProducer(protocol);
     this.resourceId = resourceId;
   }
 
@@ -67,15 +67,15 @@ public class StandaloneSpOutputCollector<T extends 
TransportProtocol> extends
 
   @Override
   public void connect() throws SpRuntimeException {
-    if (!protocolDefinition.getProducer().isConnected()) {
-      protocolDefinition.getProducer().connect(transportProtocol);
+    if (!producer.isConnected()) {
+      producer.connect();
     }
   }
 
   @Override
   public void disconnect() throws SpRuntimeException {
-    if (protocolDefinition.getProducer().isConnected()) {
-      protocolDefinition.getProducer().disconnect();
+    if (producer.isConnected()) {
+      producer.disconnect();
       ProtocolManager.removeOutputCollector(transportProtocol);
     }
   }

Reply via email to