This is an automated email from the ASF dual-hosted git repository.
riemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new 2a7ac3fb2 improvement: Remove broker dependencies from extensions
management mo… (#2117)
2a7ac3fb2 is described below
commit 2a7ac3fb29741dccddef1376b760189af872858e
Author: Dominik Riemer <[email protected]>
AuthorDate: Thu Nov 2 15:49:08 2023 +0100
improvement: Remove broker dependencies from extensions management mo…
(#2117)
* improvement: Remove broker dependencies from extensions management module
* Add pulsar dependency to extensions-all-iiot
---
streampipes-extensions-management/pom.xml | 25 ------------
.../connect/adapter/AdapterPipelineGenerator.java | 39 +------------------
.../elements/SendToBrokerAdapterSink.java | 44 ++++++++++++----------
.../elements/SendToJmsAdapterSink.java | 42 ---------------------
.../elements/SendToKafkaAdapterSink.java | 43 ---------------------
.../elements/SendToMqttAdapterSink.java | 42 ---------------------
.../elements/SendToNatsAdapterSink.java | 43 ---------------------
.../elements/SendToPulsarAdapterSink.java | 43 ---------------------
.../opcua/config/SpOpcUaConfigExtractor.java | 2 +-
.../streampipes-extensions-all-iiot/pom.xml | 5 +++
.../streampipes-extensions-all-jvm/pom.xml | 5 +++
11 files changed, 37 insertions(+), 296 deletions(-)
diff --git a/streampipes-extensions-management/pom.xml
b/streampipes-extensions-management/pom.xml
index a63d88a2f..7ce180de6 100644
--- a/streampipes-extensions-management/pom.xml
+++ b/streampipes-extensions-management/pom.xml
@@ -63,31 +63,6 @@
<artifactId>streampipes-measurement-units</artifactId>
<version>0.93.0-SNAPSHOT</version>
</dependency>
- <dependency>
- <groupId>org.apache.streampipes</groupId>
- <artifactId>streampipes-messaging-kafka</artifactId>
- <version>0.93.0-SNAPSHOT</version>
- </dependency>
- <dependency>
- <groupId>org.apache.streampipes</groupId>
- <artifactId>streampipes-messaging-jms</artifactId>
- <version>0.93.0-SNAPSHOT</version>
- </dependency>
- <dependency>
- <groupId>org.apache.streampipes</groupId>
- <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-messaging-pulsar</artifactId>
- <version>0.93.0-SNAPSHOT</version>
- </dependency>
<dependency>
<groupId>org.apache.streampipes</groupId>
<artifactId>streampipes-sdk</artifactId>
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/AdapterPipelineGenerator.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/AdapterPipelineGenerator.java
index 3fe245255..2e20e276c 100644
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/AdapterPipelineGenerator.java
+++
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/AdapterPipelineGenerator.java
@@ -19,21 +19,9 @@
package org.apache.streampipes.extensions.management.connect.adapter;
import org.apache.streampipes.connect.shared.AdapterPipelineGeneratorBase;
-import
org.apache.streampipes.extensions.management.client.StreamPipesClientResolver;
import
org.apache.streampipes.extensions.management.connect.adapter.model.pipeline.AdapterPipeline;
import
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToBrokerAdapterSink;
-import
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToJmsAdapterSink;
-import
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToKafkaAdapterSink;
-import
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToMqttAdapterSink;
-import
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToNatsAdapterSink;
-import
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements.SendToPulsarAdapterSink;
-import org.apache.streampipes.model.configuration.MessagingSettings;
-import org.apache.streampipes.model.configuration.SpProtocol;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
-import org.apache.streampipes.model.grounding.JmsTransportProtocol;
-import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
-import org.apache.streampipes.model.grounding.MqttTransportProtocol;
-import org.apache.streampipes.model.grounding.PulsarTransportProtocol;
public class AdapterPipelineGenerator extends AdapterPipelineGeneratorBase {
@@ -51,31 +39,8 @@ public class AdapterPipelineGenerator extends
AdapterPipelineGeneratorBase {
}
}
- private SendToBrokerAdapterSink<?> getAdapterSink(AdapterDescription
adapterDescription) {
- var prioritizedProtocol =
- getMessagingSettings().getPrioritizedProtocols().get(0);
-
- if (isPrioritized(prioritizedProtocol, JmsTransportProtocol.class)) {
- return new SendToJmsAdapterSink(adapterDescription);
- } else if (isPrioritized(prioritizedProtocol,
KafkaTransportProtocol.class)) {
- return new SendToKafkaAdapterSink(adapterDescription);
- } else if (isPrioritized(prioritizedProtocol,
MqttTransportProtocol.class)) {
- return new SendToMqttAdapterSink(adapterDescription);
- } else if (isPrioritized(prioritizedProtocol,
PulsarTransportProtocol.class)) {
- return new SendToPulsarAdapterSink(adapterDescription);
- } else {
- return new SendToNatsAdapterSink(adapterDescription);
- }
- }
-
- private boolean isPrioritized(SpProtocol prioritizedProtocol,
- Class<?> protocolClass) {
- return
prioritizedProtocol.getProtocolClass().equals(protocolClass.getCanonicalName());
- }
-
- private MessagingSettings getMessagingSettings() {
- var client = new
StreamPipesClientResolver().makeStreamPipesClientInstance();
- return client.adminApi().getMessagingSettings();
+ private SendToBrokerAdapterSink getAdapterSink(AdapterDescription
adapterDescription) {
+ return new SendToBrokerAdapterSink(adapterDescription);
}
private boolean hasValidGrounding(AdapterDescription adapterDescription) {
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 b9cebe258..66983df35 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
@@ -19,51 +19,52 @@ package
org.apache.streampipes.extensions.management.connect.adapter.preprocessi
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.dataformat.SpDataFormatDefinition;
import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement;
import org.apache.streampipes.extensions.api.monitoring.SpMonitoringManager;
import
org.apache.streampipes.extensions.management.connect.adapter.util.TransportFormatSelector;
import
org.apache.streampipes.extensions.management.monitoring.ExtensionsLogger;
import org.apache.streampipes.messaging.EventProducer;
+import org.apache.streampipes.messaging.SpProtocolManager;
import org.apache.streampipes.model.connect.adapter.AdapterDescription;
+import org.apache.streampipes.model.grounding.KafkaTransportProtocol;
import org.apache.streampipes.model.grounding.TransportFormat;
import org.apache.streampipes.model.grounding.TransportProtocol;
import java.util.Map;
-public abstract class SendToBrokerAdapterSink<T extends TransportProtocol>
implements IAdapterPipelineElement {
+public class SendToBrokerAdapterSink implements IAdapterPipelineElement {
protected AdapterDescription adapterDescription;
protected SpDataFormatDefinition dataFormatDefinition;
- protected T protocol;
+ protected TransportProtocol protocol;
private final EventProducer producer;
- public SendToBrokerAdapterSink(AdapterDescription adapterDescription,
- Class<T> protocolClass) {
+ public SendToBrokerAdapterSink(AdapterDescription adapterDescription) {
this.adapterDescription = adapterDescription;
- this.protocol = protocolClass.cast(adapterDescription
+ this.protocol = adapterDescription
.getEventGrounding()
- .getTransportProtocol());
+ .getTransportProtocol();
if (getEnvironment().getSpDebug().getValueOrDefault()) {
modifyProtocolForDebugging(this.protocol);
}
- this.producer = makeProducer(this.protocol);
+ var producerOpt = SpProtocolManager.INSTANCE.findDefinition(this.protocol);
+ if (producerOpt.isPresent()) {
+ this.producer = producerOpt.get().getProducer(this.protocol);
- TransportFormat transportFormat = adapterDescription
- .getEventGrounding()
- .getTransportFormats()
- .get(0);
+ TransportFormat transportFormat = adapterDescription
+ .getEventGrounding()
+ .getTransportFormats()
+ .get(0);
- this.dataFormatDefinition =
- new TransportFormatSelector(transportFormat).getDataFormatDefinition();
+ this.dataFormatDefinition =
+ new
TransportFormatSelector(transportFormat).getDataFormatDefinition();
- try {
producer.connect();
- } catch (SpRuntimeException e) {
- e.printStackTrace();
+ } else {
+ throw new RuntimeException("Could not find protocol");
}
}
@@ -86,9 +87,12 @@ public abstract class SendToBrokerAdapterSink<T extends
TransportProtocol> imple
producer.publish(event);
}
- protected abstract EventProducer makeProducer(T protocol);
-
- public abstract void modifyProtocolForDebugging(T protocol);
+ public void modifyProtocolForDebugging(TransportProtocol protocol) {
+ protocol.setBrokerHostname("localhost");
+ if (protocol instanceof KafkaTransportProtocol) {
+ ((KafkaTransportProtocol) protocol).setKafkaPort(9094);
+ }
+ }
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
deleted file mode 100644
index 98af0b42c..000000000
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToJmsAdapterSink.java
+++ /dev/null
@@ -1,42 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- *
- */
-package
org.apache.streampipes.extensions.management.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;
-
-public class SendToJmsAdapterSink extends
SendToBrokerAdapterSink<JmsTransportProtocol>
- implements IAdapterPipelineElement {
-
- public SendToJmsAdapterSink(AdapterDescription adapterDescription) {
- super(adapterDescription, JmsTransportProtocol.class);
- }
-
- @Override
- protected EventProducer makeProducer(JmsTransportProtocol protocol) {
- return new ActiveMQPublisher(protocol);
- }
-
- @Override
- public void modifyProtocolForDebugging(JmsTransportProtocol protocol) {
- protocol.setBrokerHostname("localhost");
- }
-}
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
deleted file mode 100644
index 5db6d32b2..000000000
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToKafkaAdapterSink.java
+++ /dev/null
@@ -1,43 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- *
- */
-package
org.apache.streampipes.extensions.management.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;
-
-public class SendToKafkaAdapterSink extends
SendToBrokerAdapterSink<KafkaTransportProtocol>
- implements IAdapterPipelineElement {
-
- public SendToKafkaAdapterSink(AdapterDescription adapterDescription) {
- super(adapterDescription, KafkaTransportProtocol.class);
- }
-
- @Override
- protected EventProducer makeProducer(KafkaTransportProtocol protocol) {
- return new SpKafkaProducer(protocol);
- }
-
- @Override
- public void modifyProtocolForDebugging(KafkaTransportProtocol protocol) {
- protocol.setBrokerHostname("localhost");
- protocol.setKafkaPort(9094);
- }
-}
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
deleted file mode 100644
index d2823c1c5..000000000
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToMqttAdapterSink.java
+++ /dev/null
@@ -1,42 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- *
- */
-package
org.apache.streampipes.extensions.management.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;
-
-public class SendToMqttAdapterSink extends
SendToBrokerAdapterSink<MqttTransportProtocol>
- implements IAdapterPipelineElement {
-
- public SendToMqttAdapterSink(AdapterDescription adapterDescription) {
- super(adapterDescription, MqttTransportProtocol.class);
- }
-
- @Override
- protected EventProducer makeProducer(MqttTransportProtocol protocol) {
- return new MqttPublisher(protocol);
- }
-
- @Override
- public void modifyProtocolForDebugging(MqttTransportProtocol
transportProtocol) {
- protocol.setBrokerHostname("localhost");
- }
-}
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
deleted file mode 100644
index 5216ee0b0..000000000
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToNatsAdapterSink.java
+++ /dev/null
@@ -1,43 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- *
- */
-
-package
org.apache.streampipes.extensions.management.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;
-
-public class SendToNatsAdapterSink extends
SendToBrokerAdapterSink<NatsTransportProtocol>
- implements IAdapterPipelineElement {
-
- public SendToNatsAdapterSink(AdapterDescription adapterDescription) {
- super(adapterDescription, NatsTransportProtocol.class);
- }
-
- @Override
- protected EventProducer makeProducer(NatsTransportProtocol protocol) {
- return new NatsPublisher(protocol);
- }
-
- @Override
- public void modifyProtocolForDebugging(NatsTransportProtocol protocol) {
- protocol.setBrokerHostname("localhost");
- }
-}
diff --git
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToPulsarAdapterSink.java
b/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToPulsarAdapterSink.java
deleted file mode 100644
index b9f3ba56d..000000000
---
a/streampipes-extensions-management/src/main/java/org/apache/streampipes/extensions/management/connect/adapter/preprocessing/elements/SendToPulsarAdapterSink.java
+++ /dev/null
@@ -1,43 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- *
- */
-
-package
org.apache.streampipes.extensions.management.connect.adapter.preprocessing.elements;
-
-import org.apache.streampipes.extensions.api.connect.IAdapterPipelineElement;
-import org.apache.streampipes.messaging.EventProducer;
-import org.apache.streampipes.messaging.pulsar.PulsarProducer;
-import org.apache.streampipes.model.connect.adapter.AdapterDescription;
-import org.apache.streampipes.model.grounding.PulsarTransportProtocol;
-
-public class SendToPulsarAdapterSink extends
SendToBrokerAdapterSink<PulsarTransportProtocol>
- implements IAdapterPipelineElement {
-
- public SendToPulsarAdapterSink(AdapterDescription adapterDescription) {
- super(adapterDescription, PulsarTransportProtocol.class);
- }
-
- @Override
- protected EventProducer makeProducer(PulsarTransportProtocol protocol) {
- return new PulsarProducer(protocol);
- }
-
- @Override
- public void modifyProtocolForDebugging(PulsarTransportProtocol protocol) {
- protocol.setBrokerHostname("pulsar://localhost:6650");
- }
-}
diff --git
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java
index 2c3e3b138..612b3e490 100644
---
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java
+++
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java
@@ -24,7 +24,6 @@ import
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaUtil;
import java.util.List;
-import static org.apache.kafka.common.config.ConfigDef.Type.PASSWORD;
import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.ACCESS_MODE;
import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.ADAPTER_TYPE;
import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.AVAILABLE_NODES;
@@ -33,6 +32,7 @@ import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabe
import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.OPC_SERVER_PORT;
import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.OPC_SERVER_URL;
import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.OPC_URL;
+import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PASSWORD;
import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PULLING_INTERVAL;
import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PULL_MODE;
import static
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.UNAUTHENTICATED;
diff --git a/streampipes-extensions/streampipes-extensions-all-iiot/pom.xml
b/streampipes-extensions/streampipes-extensions-all-iiot/pom.xml
index e47a0ab0a..b1caa7a5b 100644
--- a/streampipes-extensions/streampipes-extensions-all-iiot/pom.xml
+++ b/streampipes-extensions/streampipes-extensions-all-iiot/pom.xml
@@ -61,6 +61,11 @@
<artifactId>streampipes-connectors-influx</artifactId>
<version>0.93.0-SNAPSHOT</version>
</dependency>
+ <dependency>
+ <groupId>org.apache.streampipes</groupId>
+ <artifactId>streampipes-messaging-pulsar</artifactId>
+ <version>0.93.0-SNAPSHOT</version>
+ </dependency>
<dependency>
<groupId>org.apache.streampipes</groupId>
<artifactId>streampipes-processors-filters-siddhi</artifactId>
diff --git a/streampipes-extensions/streampipes-extensions-all-jvm/pom.xml
b/streampipes-extensions/streampipes-extensions-all-jvm/pom.xml
index 3b03c5940..bfaa6dc03 100644
--- a/streampipes-extensions/streampipes-extensions-all-jvm/pom.xml
+++ b/streampipes-extensions/streampipes-extensions-all-jvm/pom.xml
@@ -64,6 +64,11 @@
<classifier>embed</classifier>
</dependency>
+ <dependency>
+ <groupId>org.apache.streampipes</groupId>
+ <artifactId>streampipes-messaging-pulsar</artifactId>
+ <version>0.93.0-SNAPSHOT</version>
+ </dependency>
<dependency>
<groupId>org.apache.streampipes</groupId>
<artifactId>streampipes-messaging-nats</artifactId>