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