From a6e090ef868fab068c21a65dc76a11c8c240e780 Mon Sep 17 00:00:00 2001
From: Yevhen Bondarenko <56396344+YevhenBondarenko@users.noreply.github.com>
Date: Tue, 31 Mar 2020 16:39:41 +0300
Subject: [PATCH] Develop/2.5 pubsub (#2566)
* created Aws Sqs Queue
* improvement AwsSqs providers
* created pubsub queue
* revert package-lock.json
* Aws sqs improvements
* Aws sqs improvements
* Aws sqs improvements
* Aws sqs improvements
* Created pubsub queue
* aws improvements
* aws improvements
* aws improvements
* added visibility timeout to aws queue
* pub sub improvements
* pub sub improvements
* aws sqs improvements
* pub sub improvements
* added comment to transport.yml about ack deadline
---
.../src/main/resources/thingsboard.yml | 10 +-
common/queue/pom.xml | 5 +-
.../DefaultTbQueueMsg.java} | 7 +-
.../common/DefaultTbQueueRequestTemplate.java | 2 +-
.../provider/KafkaTbCoreQueueProvider.java | 6 +-
.../KafkaTbRuleEngineQueueProvider.java | 10 +-
.../provider/PubSubMonolithQueueProvider.java | 137 +++++++++++
.../provider/PubSubTbCoreQueueProvider.java | 113 +++++++++
.../PubSubTbRuleEngineQueueProvider.java | 101 ++++++++
.../PubSubTransportQueueProvider.java | 104 ++++++++
.../server/queue/pubsub/TbPubSubAdmin.java | 157 ++++++++++++
.../pubsub/TbPubSubConsumerTemplate.java | 228 ++++++++++++++++++
.../pubsub/TbPubSubProducerTemplate.java | 135 +++++++++++
.../server/queue/pubsub/TbPubSubSettings.java | 61 +++++
.../queue/sqs/TbAwsSqsConsumerTemplate.java | 3 +-
.../queue/sqs/TbAwsSqsProducerTemplate.java | 3 +-
pom.xml | 16 +-
rule-engine/rule-engine-components/pom.xml | 2 -
.../src/main/resources/tb-mqtt-transport.yml | 10 +-
19 files changed, 1078 insertions(+), 32 deletions(-)
rename common/queue/src/main/java/org/thingsboard/server/queue/{sqs/TbAwsSqsMsg.java => common/DefaultTbQueueMsg.java} (86%)
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueProvider.java
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueProvider.java
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueProvider.java
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueProvider.java
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubProducerTemplate.java
create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubSettings.java
diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml
index 2fe8f59608..e731bd7bac 100644
--- a/application/src/main/resources/thingsboard.yml
+++ b/application/src/main/resources/thingsboard.yml
@@ -517,7 +517,7 @@ swagger:
version: "${SWAGGER_VERSION:2.0}"
queue:
- type: "${TB_QUEUE_TYPE:in-memory}" # kafka or in-memory or aws-sqs
+ type: "${TB_QUEUE_TYPE:in-memory}" # kafka or in-memory or aws-sqs or pubsub
kafka:
bootstrap.servers: "${TB_KAFKA_SERVERS:localhost:9092}"
acks: "${TB_KAFKA_ACKS:all}"
@@ -530,7 +530,13 @@ queue:
secret_access_key: "${TB_QUEUE_AWS_SQS_SECRET_ACCESS_KEY:YOUR_SECRET}"
region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}"
threads_per_topic: "${TB_QUEUE_AWS_SQS_THREADS_PER_TOPIC:1}"
- visibility_timeout: "${TB_QUEUE_AWS_SQS_VISIBILITY_TIMEOUT:30}" #in seconds
+ visibility_timeout: "${TB_QUEUE_AWS_SQS_VISIBILITY_TIMEOUT:30}" #In seconds. If messages wont commit in this time, messages will poll again
+ pubsub:
+ project_id: "${TB_QUEUE_PUBSUB_PROJECT_ID:YOUR_PROJECT_ID}"
+ service_account: "${TB_QUEUE_PUBSUB_SERVICE_ACCOUNT:YOUR_SERVICE_ACCOUNT}"
+ ack_deadline: "${TB_QUEUE_PUBSUB_ACK_DEADLINE:30}" #In seconds. If messages wont commit in this time, messages will poll again
+ max_msg_size: "${TB_QUEUE_PUBSUB_MAX_MSG_SIZE:1048576}" #in bytes
+ max_messages: "${TB_QUEUE_PUBSUB_MAX_MESSAGES:1000}"
partitions:
hash_function_name: "${TB_QUEUE_PARTITIONS_HASH_FUNCTION_NAME:murmur3_128}"
virtual_nodes_size: "${TB_QUEUE_PARTITIONS_VIRTUAL_NODES_SIZE:16}"
diff --git a/common/queue/pom.xml b/common/queue/pom.xml
index e201b8995f..0a61eb0cf6 100644
--- a/common/queue/pom.xml
+++ b/common/queue/pom.xml
@@ -56,7 +56,10 @@
com.amazonaws
aws-java-sdk-sqs
-
+
+ com.google.cloud
+ google-cloud-pubsub
+
org.springframework
spring-context-support
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsMsg.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueMsg.java
similarity index 86%
rename from common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsMsg.java
rename to common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueMsg.java
index 4df7558491..0e816ae59b 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsMsg.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueMsg.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.queue.sqs;
+package org.thingsboard.server.queue.common;
import com.google.gson.annotations.Expose;
import lombok.Data;
@@ -23,16 +23,15 @@ import org.thingsboard.server.queue.TbQueueMsgHeaders;
import java.util.UUID;
@Data
-public class TbAwsSqsMsg implements TbQueueMsg {
+public class DefaultTbQueueMsg implements TbQueueMsg {
private final UUID key;
private final byte[] data;
- public TbAwsSqsMsg(UUID key, byte[] data) {
+ public DefaultTbQueueMsg(UUID key, byte[] data) {
this.key = key;
this.data = data;
}
@Expose(serialize = false, deserialize = false)
private TbQueueMsgHeaders headers;
-
}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java
index 40ca2219dd..6ae9b9a33c 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java
@@ -96,7 +96,7 @@ public class DefaultTbQueueRequestTemplate {
- log.trace("Received response to Kafka Template request: {}", response);
+ log.trace("Received response to Queue Template request: {}", response);
byte[] requestIdHeader = response.getHeaders().get(REQUEST_ID_HEADER);
UUID requestId;
if (requestIdHeader == null) {
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueProvider.java
index 177f5f3c3d..110a98b3b9 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueProvider.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueProvider.java
@@ -30,7 +30,6 @@ import org.thingsboard.server.queue.TbQueueCoreSettings;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.TbQueueRuleEngineSettings;
import org.thingsboard.server.queue.TbQueueTransportApiSettings;
-import org.thingsboard.server.queue.TbQueueTransportNotificationSettings;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
@@ -48,21 +47,18 @@ public class KafkaTbCoreQueueProvider implements TbCoreQueueProvider {
private final TbQueueCoreSettings coreSettings;
private final TbQueueRuleEngineSettings ruleEngineSettings;
private final TbQueueTransportApiSettings transportApiSettings;
- private final TbQueueTransportNotificationSettings transportNotificationSettings;
public KafkaTbCoreQueueProvider(PartitionService partitionService, TbKafkaSettings kafkaSettings,
TbServiceInfoProvider serviceInfoProvider,
TbQueueCoreSettings coreSettings,
TbQueueRuleEngineSettings ruleEngineSettings,
- TbQueueTransportApiSettings transportApiSettings,
- TbQueueTransportNotificationSettings transportNotificationSettings) {
+ TbQueueTransportApiSettings transportApiSettings) {
this.partitionService = partitionService;
this.kafkaSettings = kafkaSettings;
this.serviceInfoProvider = serviceInfoProvider;
this.coreSettings = coreSettings;
this.ruleEngineSettings = ruleEngineSettings;
this.transportApiSettings = transportApiSettings;
- this.transportNotificationSettings = transportNotificationSettings;
}
@Override
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueProvider.java
index 36ffd05b48..7f939ee11d 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueProvider.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbRuleEngineQueueProvider.java
@@ -27,8 +27,6 @@ import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.TbQueueCoreSettings;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.TbQueueRuleEngineSettings;
-import org.thingsboard.server.queue.TbQueueTransportApiSettings;
-import org.thingsboard.server.queue.TbQueueTransportNotificationSettings;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
@@ -45,22 +43,16 @@ public class KafkaTbRuleEngineQueueProvider implements TbRuleEngineQueueProvider
private final TbServiceInfoProvider serviceInfoProvider;
private final TbQueueCoreSettings coreSettings;
private final TbQueueRuleEngineSettings ruleEngineSettings;
- private final TbQueueTransportApiSettings transportApiSettings;
- private final TbQueueTransportNotificationSettings transportNotificationSettings;
public KafkaTbRuleEngineQueueProvider(PartitionService partitionService, TbKafkaSettings kafkaSettings,
TbServiceInfoProvider serviceInfoProvider,
TbQueueCoreSettings coreSettings,
- TbQueueRuleEngineSettings ruleEngineSettings,
- TbQueueTransportApiSettings transportApiSettings,
- TbQueueTransportNotificationSettings transportNotificationSettings) {
+ TbQueueRuleEngineSettings ruleEngineSettings) {
this.partitionService = partitionService;
this.kafkaSettings = kafkaSettings;
this.serviceInfoProvider = serviceInfoProvider;
this.coreSettings = coreSettings;
this.ruleEngineSettings = ruleEngineSettings;
- this.transportApiSettings = transportApiSettings;
- this.transportNotificationSettings = transportNotificationSettings;
}
@Override
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueProvider.java
new file mode 100644
index 0000000000..2dc9549679
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubMonolithQueueProvider.java
@@ -0,0 +1,137 @@
+/**
+ * Copyright © 2016-2020 The Thingsboard Authors
+ *
+ * Licensed 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.thingsboard.server.queue.provider;
+
+import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
+import org.springframework.stereotype.Component;
+import org.thingsboard.server.common.msg.queue.ServiceType;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
+import org.thingsboard.server.queue.TbQueueAdmin;
+import org.thingsboard.server.queue.TbQueueConsumer;
+import org.thingsboard.server.queue.TbQueueCoreSettings;
+import org.thingsboard.server.queue.TbQueueProducer;
+import org.thingsboard.server.queue.TbQueueRuleEngineSettings;
+import org.thingsboard.server.queue.TbQueueTransportApiSettings;
+import org.thingsboard.server.queue.TbQueueTransportNotificationSettings;
+import org.thingsboard.server.queue.common.TbProtoQueueMsg;
+import org.thingsboard.server.queue.discovery.PartitionService;
+import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
+import org.thingsboard.server.queue.pubsub.TbPubSubAdmin;
+import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate;
+import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate;
+import org.thingsboard.server.queue.pubsub.TbPubSubSettings;
+
+@Component
+@ConditionalOnExpression("'${queue.type:null}'=='pubsub' && '${service.type:null}'=='monolith'")
+public class PubSubMonolithQueueProvider implements TbCoreQueueProvider, TbRuleEngineQueueProvider {
+
+ private final TbPubSubSettings pubSubSettings;
+ private final TbQueueCoreSettings coreSettings;
+ private final TbQueueRuleEngineSettings ruleEngineSettings;
+ private final TbQueueTransportApiSettings transportApiSettings;
+ private final TbQueueTransportNotificationSettings transportNotificationSettings;
+ private final TbQueueAdmin admin;
+ private final PartitionService partitionService;
+ private final TbServiceInfoProvider serviceInfoProvider;
+
+ private TbQueueProducer> tbCoreProducer;
+
+ public PubSubMonolithQueueProvider(TbPubSubSettings pubSubSettings,
+ TbQueueCoreSettings coreSettings,
+ TbQueueRuleEngineSettings ruleEngineSettings,
+ TbQueueTransportApiSettings transportApiSettings,
+ TbQueueTransportNotificationSettings transportNotificationSettings,
+ PartitionService partitionService,
+ TbServiceInfoProvider serviceInfoProvider) {
+ this.pubSubSettings = pubSubSettings;
+ this.coreSettings = coreSettings;
+ this.ruleEngineSettings = ruleEngineSettings;
+ this.transportApiSettings = transportApiSettings;
+ this.transportNotificationSettings = transportNotificationSettings;
+ this.admin = new TbPubSubAdmin(pubSubSettings);
+ this.partitionService = partitionService;
+ this.serviceInfoProvider = serviceInfoProvider;
+ }
+
+ @Override
+ public TbQueueProducer> getTransportNotificationsMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, transportNotificationSettings.getNotificationsTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getRuleEngineMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, ruleEngineSettings.getTopic());
+
+ }
+
+ @Override
+ public TbQueueProducer> getRuleEngineNotificationsMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, ruleEngineSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getTbCoreMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getTbCoreNotificationsMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueConsumer> getToRuleEngineMsgConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings, ruleEngineSettings.getTopic(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+
+ @Override
+ public TbQueueConsumer> getToRuleEngineNotificationsMsgConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings,
+ partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+
+ @Override
+ public TbQueueConsumer> getToCoreMsgConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings, coreSettings.getTopic(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+
+ @Override
+ public TbQueueConsumer> getToCoreNotificationsMsgConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings,
+ partitionService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+
+ @Override
+ public TbQueueConsumer> getTransportApiRequestConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings, transportApiSettings.getRequestsTopic(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+
+ @Override
+ public TbQueueProducer> getTransportApiResponseProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, transportApiSettings.getResponsesTopic());
+ }
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueProvider.java
new file mode 100644
index 0000000000..4c89b35dec
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbCoreQueueProvider.java
@@ -0,0 +1,113 @@
+/**
+ * Copyright © 2016-2020 The Thingsboard Authors
+ *
+ * Licensed 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.thingsboard.server.queue.provider;
+
+import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
+import org.springframework.stereotype.Component;
+import org.thingsboard.server.common.msg.queue.ServiceType;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
+import org.thingsboard.server.queue.TbQueueAdmin;
+import org.thingsboard.server.queue.TbQueueConsumer;
+import org.thingsboard.server.queue.TbQueueCoreSettings;
+import org.thingsboard.server.queue.TbQueueProducer;
+import org.thingsboard.server.queue.TbQueueTransportApiSettings;
+import org.thingsboard.server.queue.common.TbProtoQueueMsg;
+import org.thingsboard.server.queue.discovery.PartitionService;
+import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
+import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate;
+import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate;
+import org.thingsboard.server.queue.pubsub.TbPubSubSettings;
+
+@Component
+@ConditionalOnExpression("'${queue.type:null}'=='pubsub' && '${service.type:null}'=='tb-core'")
+public class PubSubTbCoreQueueProvider implements TbCoreQueueProvider {
+
+ private final TbPubSubSettings pubSubSettings;
+ private final TbQueueCoreSettings coreSettings;
+ private final TbQueueTransportApiSettings transportApiSettings;
+ private final TbQueueAdmin admin;
+ private final PartitionService partitionService;
+ private final TbServiceInfoProvider serviceInfoProvider;
+
+ public PubSubTbCoreQueueProvider(TbPubSubSettings pubSubSettings,
+ TbQueueCoreSettings coreSettings,
+ TbQueueTransportApiSettings transportApiSettings,
+ TbQueueAdmin admin,
+ PartitionService partitionService,
+ TbServiceInfoProvider serviceInfoProvider) {
+ this.pubSubSettings = pubSubSettings;
+ this.coreSettings = coreSettings;
+ this.transportApiSettings = transportApiSettings;
+ this.admin = admin;
+ this.partitionService = partitionService;
+ this.serviceInfoProvider = serviceInfoProvider;
+ }
+
+ @Override
+ public TbQueueProducer> getTransportNotificationsMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getRuleEngineMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getRuleEngineNotificationsMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getTbCoreMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getTbCoreNotificationsMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueConsumer> getToCoreMsgConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings, coreSettings.getTopic(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+
+ @Override
+ public TbQueueConsumer> getToCoreNotificationsMsgConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings,
+ partitionService.getNotificationsTopic(ServiceType.TB_CORE, serviceInfoProvider.getServiceId()).getFullTopicName(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), ToCoreNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+
+ @Override
+ public TbQueueConsumer> getTransportApiRequestConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings, transportApiSettings.getRequestsTopic(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiRequestMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+
+ @Override
+ public TbQueueProducer> getTransportApiResponseProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueProvider.java
new file mode 100644
index 0000000000..3f707235fa
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTbRuleEngineQueueProvider.java
@@ -0,0 +1,101 @@
+/**
+ * Copyright © 2016-2020 The Thingsboard Authors
+ *
+ * Licensed 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.thingsboard.server.queue.provider;
+
+import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
+import org.springframework.stereotype.Component;
+import org.thingsboard.server.common.msg.queue.ServiceType;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
+import org.thingsboard.server.queue.TbQueueAdmin;
+import org.thingsboard.server.queue.TbQueueConsumer;
+import org.thingsboard.server.queue.TbQueueCoreSettings;
+import org.thingsboard.server.queue.TbQueueProducer;
+import org.thingsboard.server.queue.TbQueueRuleEngineSettings;
+import org.thingsboard.server.queue.common.TbProtoQueueMsg;
+import org.thingsboard.server.queue.discovery.PartitionService;
+import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
+import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate;
+import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate;
+import org.thingsboard.server.queue.pubsub.TbPubSubSettings;
+
+@Component
+@ConditionalOnExpression("'${queue.type:null}'=='pubsub' && '${service.type:null}'=='tb-rule-engine'")
+public class PubSubTbRuleEngineQueueProvider implements TbRuleEngineQueueProvider {
+
+ private final TbPubSubSettings pubSubSettings;
+ private final TbQueueCoreSettings coreSettings;
+ private final TbQueueRuleEngineSettings ruleEngineSettings;
+ private final TbQueueAdmin admin;
+ private final PartitionService partitionService;
+ private final TbServiceInfoProvider serviceInfoProvider;
+
+ public PubSubTbRuleEngineQueueProvider(TbPubSubSettings pubSubSettings,
+ TbQueueCoreSettings coreSettings,
+ TbQueueRuleEngineSettings ruleEngineSettings,
+ TbQueueAdmin admin,
+ PartitionService partitionService,
+ TbServiceInfoProvider serviceInfoProvider) {
+ this.pubSubSettings = pubSubSettings;
+ this.coreSettings = coreSettings;
+ this.ruleEngineSettings = ruleEngineSettings;
+ this.admin = admin;
+ this.partitionService = partitionService;
+ this.serviceInfoProvider = serviceInfoProvider;
+ }
+
+ @Override
+ public TbQueueProducer> getTransportNotificationsMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getRuleEngineMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getRuleEngineNotificationsMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, ruleEngineSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getTbCoreMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+
+ }
+
+ @Override
+ public TbQueueProducer> getTbCoreNotificationsMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueConsumer> getToRuleEngineMsgConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings, ruleEngineSettings.getTopic(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+
+ @Override
+ public TbQueueConsumer> getToRuleEngineNotificationsMsgConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings,
+ partitionService.getNotificationsTopic(ServiceType.TB_RULE_ENGINE, serviceInfoProvider.getServiceId()).getFullTopicName(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), ToRuleEngineNotificationMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueProvider.java
new file mode 100644
index 0000000000..b1a7efc950
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/PubSubTransportQueueProvider.java
@@ -0,0 +1,104 @@
+/**
+ * Copyright © 2016-2020 The Thingsboard Authors
+ *
+ * Licensed 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.thingsboard.server.queue.provider;
+
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
+import org.springframework.stereotype.Component;
+import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
+import org.thingsboard.server.queue.TbQueueAdmin;
+import org.thingsboard.server.queue.TbQueueConsumer;
+import org.thingsboard.server.queue.TbQueueCoreSettings;
+import org.thingsboard.server.queue.TbQueueProducer;
+import org.thingsboard.server.queue.TbQueueRequestTemplate;
+import org.thingsboard.server.queue.TbQueueRuleEngineSettings;
+import org.thingsboard.server.queue.TbQueueTransportApiSettings;
+import org.thingsboard.server.queue.TbQueueTransportNotificationSettings;
+import org.thingsboard.server.queue.common.DefaultTbQueueRequestTemplate;
+import org.thingsboard.server.queue.common.TbProtoQueueMsg;
+import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
+import org.thingsboard.server.queue.pubsub.TbPubSubAdmin;
+import org.thingsboard.server.queue.pubsub.TbPubSubConsumerTemplate;
+import org.thingsboard.server.queue.pubsub.TbPubSubProducerTemplate;
+import org.thingsboard.server.queue.pubsub.TbPubSubSettings;
+
+@Component
+@ConditionalOnExpression("'${queue.type:null}'=='pubsub' && ('${service.type:null}'=='monolith' || '${service.type:null}'=='tb-transport')")
+@Slf4j
+public class PubSubTransportQueueProvider implements TbTransportQueueProvider {
+
+ private final TbPubSubSettings pubSubSettings;
+ private final TbServiceInfoProvider serviceInfoProvider;
+ private final TbQueueCoreSettings coreSettings;
+ private final TbQueueRuleEngineSettings ruleEngineSettings;
+ private final TbQueueTransportApiSettings transportApiSettings;
+ private final TbQueueTransportNotificationSettings transportNotificationSettings;
+ private final TbQueueAdmin admin;
+
+ public PubSubTransportQueueProvider(TbPubSubSettings pubSubSettings,
+ TbServiceInfoProvider serviceInfoProvider,
+ TbQueueCoreSettings coreSettings,
+ TbQueueRuleEngineSettings ruleEngineSettings,
+ TbQueueTransportApiSettings transportApiSettings,
+ TbQueueTransportNotificationSettings transportNotificationSettings) {
+ this.pubSubSettings = pubSubSettings;
+ this.serviceInfoProvider = serviceInfoProvider;
+ this.coreSettings = coreSettings;
+ this.ruleEngineSettings = ruleEngineSettings;
+ this.transportApiSettings = transportApiSettings;
+ this.transportNotificationSettings = transportNotificationSettings;
+ this.admin = new TbPubSubAdmin(pubSubSettings);
+ }
+
+ @Override
+ public TbQueueRequestTemplate, TbProtoQueueMsg> getTransportApiRequestTemplate() {
+ TbQueueProducer> producer = new TbPubSubProducerTemplate<>(admin, pubSubSettings, transportApiSettings.getRequestsTopic());
+ TbQueueConsumer> consumer = new TbPubSubConsumerTemplate<>(admin, pubSubSettings,
+ transportApiSettings.getResponsesTopic() + "." + serviceInfoProvider.getServiceId(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), TransportApiResponseMsg.parseFrom(msg.getData()), msg.getHeaders()));
+
+ DefaultTbQueueRequestTemplate.DefaultTbQueueRequestTemplateBuilder
+ , TbProtoQueueMsg> templateBuilder = DefaultTbQueueRequestTemplate.builder();
+ templateBuilder.queueAdmin(admin);
+ templateBuilder.requestTemplate(producer);
+ templateBuilder.responseTemplate(consumer);
+ templateBuilder.maxPendingRequests(transportApiSettings.getMaxPendingRequests());
+ templateBuilder.maxRequestTimeout(transportApiSettings.getMaxRequestsTimeout());
+ templateBuilder.pollInterval(transportApiSettings.getResponsePollInterval());
+ return templateBuilder.build();
+ }
+
+ @Override
+ public TbQueueProducer> getRuleEngineMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, ruleEngineSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueProducer> getTbCoreMsgProducer() {
+ return new TbPubSubProducerTemplate<>(admin, pubSubSettings, coreSettings.getTopic());
+ }
+
+ @Override
+ public TbQueueConsumer> getTransportNotificationsConsumer() {
+ return new TbPubSubConsumerTemplate<>(admin, pubSubSettings,
+ transportNotificationSettings.getNotificationsTopic() + "." + serviceInfoProvider.getServiceId(),
+ msg -> new TbProtoQueueMsg<>(msg.getKey(), ToTransportMsg.parseFrom(msg.getData()), msg.getHeaders()));
+ }
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java
new file mode 100644
index 0000000000..f0af639d4d
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java
@@ -0,0 +1,157 @@
+/**
+ * Copyright © 2016-2020 The Thingsboard Authors
+ *
+ * Licensed 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.thingsboard.server.queue.pubsub;
+
+import com.google.cloud.pubsub.v1.SubscriptionAdminClient;
+import com.google.cloud.pubsub.v1.SubscriptionAdminSettings;
+import com.google.cloud.pubsub.v1.TopicAdminClient;
+import com.google.cloud.pubsub.v1.TopicAdminSettings;
+import com.google.pubsub.v1.ListSubscriptionsRequest;
+import com.google.pubsub.v1.ListTopicsRequest;
+import com.google.pubsub.v1.ProjectName;
+import com.google.pubsub.v1.ProjectSubscriptionName;
+import com.google.pubsub.v1.ProjectTopicName;
+import com.google.pubsub.v1.PushConfig;
+import com.google.pubsub.v1.Subscription;
+import com.google.pubsub.v1.Topic;
+import lombok.extern.slf4j.Slf4j;
+import org.thingsboard.server.queue.TbQueueAdmin;
+
+import java.io.IOException;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+
+@Slf4j
+public class TbPubSubAdmin implements TbQueueAdmin {
+
+ private final TbPubSubSettings pubSubSettings;
+ private final SubscriptionAdminSettings subscriptionAdminSettings;
+ private final TopicAdminSettings topicAdminSettings;
+ private final Set topicSet = ConcurrentHashMap.newKeySet();
+ private final Set subscriptionSet = ConcurrentHashMap.newKeySet();
+
+ public TbPubSubAdmin(TbPubSubSettings pubSubSettings) {
+ this.pubSubSettings = pubSubSettings;
+
+ try {
+ topicAdminSettings = TopicAdminSettings.newBuilder().setCredentialsProvider(pubSubSettings.getCredentialsProvider()).build();
+ } catch (IOException e) {
+ log.error("Failed to create TopicAdminSettings");
+ throw new RuntimeException("Failed to create TopicAdminSettings.");
+ }
+
+ try {
+ subscriptionAdminSettings = SubscriptionAdminSettings.newBuilder().setCredentialsProvider(pubSubSettings.getCredentialsProvider()).build();
+ } catch (IOException e) {
+ log.error("Failed to create SubscriptionAdminSettings");
+ throw new RuntimeException("Failed to create SubscriptionAdminSettings.");
+ }
+
+ try (TopicAdminClient topicAdminClient = TopicAdminClient.create(topicAdminSettings)) {
+ ListTopicsRequest listTopicsRequest =
+ ListTopicsRequest.newBuilder().setProject(ProjectName.format(pubSubSettings.getProjectId())).build();
+ TopicAdminClient.ListTopicsPagedResponse response = topicAdminClient.listTopics(listTopicsRequest);
+ for (Topic topic : response.iterateAll()) {
+ topicSet.add(topic.getName());
+ }
+ } catch (IOException e) {
+ log.error("Failed to get topics.", e);
+ throw new RuntimeException("Failed to get topics.", e);
+ }
+
+ try (SubscriptionAdminClient subscriptionAdminClient = SubscriptionAdminClient.create(subscriptionAdminSettings)) {
+
+ ListSubscriptionsRequest listSubscriptionsRequest =
+ ListSubscriptionsRequest.newBuilder()
+ .setProject(ProjectName.of(pubSubSettings.getProjectId()).toString())
+ .build();
+ SubscriptionAdminClient.ListSubscriptionsPagedResponse response =
+ subscriptionAdminClient.listSubscriptions(listSubscriptionsRequest);
+
+ for (Subscription subscription : response.iterateAll()) {
+ subscriptionSet.add(subscription.getName());
+ }
+ } catch (IOException e) {
+ log.error("Failed to get subscriptions.", e);
+ throw new RuntimeException("Failed to get subscriptions.", e);
+ }
+ }
+
+ @Override
+ public void createTopicIfNotExists(String partition) {
+ ProjectTopicName topicName = ProjectTopicName.of(pubSubSettings.getProjectId(), partition);
+
+ if (topicSet.contains(topicName.toString())) {
+ createSubscriptionIfNotExists(partition, topicName);
+ return;
+ }
+
+ try (TopicAdminClient topicAdminClient = TopicAdminClient.create(topicAdminSettings)) {
+ ListTopicsRequest listTopicsRequest =
+ ListTopicsRequest.newBuilder().setProject(ProjectName.format(pubSubSettings.getProjectId())).build();
+ TopicAdminClient.ListTopicsPagedResponse response = topicAdminClient.listTopics(listTopicsRequest);
+ for (Topic topic : response.iterateAll()) {
+ if (topic.getName().contains(topicName.toString())) {
+ topicSet.add(topic.getName());
+ createSubscriptionIfNotExists(partition, topicName);
+ return;
+ }
+ }
+
+ topicAdminClient.createTopic(topicName);
+ topicSet.add(topicName.toString());
+ log.info("Created new topic: [{}]", topicName.toString());
+ createSubscriptionIfNotExists(partition, topicName);
+ } catch (IOException e) {
+ log.error("Failed to create topic: [{}].", topicName.toString(), e);
+ throw new RuntimeException("Failed to create topic.", e);
+ }
+ }
+
+ private void createSubscriptionIfNotExists(String partition, ProjectTopicName topicName) {
+ ProjectSubscriptionName subscriptionName =
+ ProjectSubscriptionName.of(pubSubSettings.getProjectId(), partition);
+
+ if (subscriptionSet.contains(subscriptionName.toString())) {
+ return;
+ }
+
+ try (SubscriptionAdminClient subscriptionAdminClient = SubscriptionAdminClient.create(subscriptionAdminSettings)) {
+ ListSubscriptionsRequest listSubscriptionsRequest =
+ ListSubscriptionsRequest.newBuilder()
+ .setProject(ProjectName.of(pubSubSettings.getProjectId()).toString())
+ .build();
+ SubscriptionAdminClient.ListSubscriptionsPagedResponse response =
+ subscriptionAdminClient.listSubscriptions(listSubscriptionsRequest);
+
+ for (Subscription subscription : response.iterateAll()) {
+ if (subscription.getName().equals(subscriptionName.toString())) {
+ subscriptionSet.add(subscription.getName());
+ return;
+ }
+ }
+
+ subscriptionAdminClient.createSubscription(
+ subscriptionName, topicName, PushConfig.getDefaultInstance(), pubSubSettings.getAckDeadline()).getName();
+ subscriptionSet.add(subscriptionName.toString());
+ log.info("Created new subscription: [{}]", subscriptionName.toString());
+ } catch (IOException e) {
+ log.error("Failed to create subscription: [{}].", subscriptionName.toString(), e);
+ throw new RuntimeException("Failed to create subscription.", e);
+ }
+ }
+
+}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java
new file mode 100644
index 0000000000..4109e72070
--- /dev/null
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubConsumerTemplate.java
@@ -0,0 +1,228 @@
+/**
+ * Copyright © 2016-2020 The Thingsboard Authors
+ *
+ * Licensed 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.thingsboard.server.queue.pubsub;
+
+import com.google.api.core.ApiFuture;
+import com.google.api.core.ApiFutures;
+import com.google.cloud.pubsub.v1.stub.GrpcSubscriberStub;
+import com.google.cloud.pubsub.v1.stub.SubscriberStub;
+import com.google.cloud.pubsub.v1.stub.SubscriberStubSettings;
+import com.google.common.reflect.TypeToken;
+import com.google.gson.Gson;
+import com.google.protobuf.InvalidProtocolBufferException;
+import com.google.pubsub.v1.AcknowledgeRequest;
+import com.google.pubsub.v1.ProjectSubscriptionName;
+import com.google.pubsub.v1.PubsubMessage;
+import com.google.pubsub.v1.PullRequest;
+import com.google.pubsub.v1.PullResponse;
+import com.google.pubsub.v1.ReceivedMessage;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.util.CollectionUtils;
+import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
+import org.thingsboard.server.queue.TbQueueAdmin;
+import org.thingsboard.server.queue.TbQueueConsumer;
+import org.thingsboard.server.queue.TbQueueMsg;
+import org.thingsboard.server.queue.TbQueueMsgDecoder;
+import org.thingsboard.server.queue.TbQueueMsgHeaders;
+import org.thingsboard.server.queue.common.DefaultTbQueueMsg;
+import org.thingsboard.server.queue.common.DefaultTbQueueMsgHeaders;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.stream.Collectors;
+
+@Slf4j
+public class TbPubSubConsumerTemplate implements TbQueueConsumer {
+
+ private final Gson gson = new Gson();
+ private final TbQueueAdmin admin;
+ private final String topic;
+ private final TbQueueMsgDecoder decoder;
+ private final TbPubSubSettings pubSubSettings;
+
+ private volatile boolean subscribed;
+ private volatile Set partitions;
+ private volatile Set subscriptionNames;
+ private final List acknowledgeRequests = new CopyOnWriteArrayList<>();
+
+ private ExecutorService consumerExecutor;
+ private final SubscriberStub subscriber;
+ private volatile boolean stopped;
+
+ private volatile int messagesPerTopic;
+
+ public TbPubSubConsumerTemplate(TbQueueAdmin admin, TbPubSubSettings pubSubSettings, String topic, TbQueueMsgDecoder decoder) {
+ this.admin = admin;
+ this.pubSubSettings = pubSubSettings;
+ this.topic = topic;
+ this.decoder = decoder;
+
+ try {
+ SubscriberStubSettings subscriberStubSettings =
+ SubscriberStubSettings.newBuilder()
+ .setCredentialsProvider(pubSubSettings.getCredentialsProvider())
+ .setTransportChannelProvider(
+ SubscriberStubSettings.defaultGrpcTransportProviderBuilder()
+ .setMaxInboundMessageSize(pubSubSettings.getMaxMsgSize())
+ .build())
+ .build();
+
+ this.subscriber = GrpcSubscriberStub.create(subscriberStubSettings);
+ } catch (IOException e) {
+ log.error("Failed to create subscriber.", e);
+ throw new RuntimeException("Failed to create subscriber.", e);
+ }
+ stopped = false;
+ }
+
+ @Override
+ public String getTopic() {
+ return topic;
+ }
+
+ @Override
+ public void subscribe() {
+ partitions = Collections.singleton(new TopicPartitionInfo(topic, null, null, true));
+ subscribed = false;
+ }
+
+ @Override
+ public void subscribe(Set partitions) {
+ this.partitions = partitions;
+ subscribed = false;
+ }
+
+ @Override
+ public void unsubscribe() {
+ stopped = true;
+ if (consumerExecutor != null) {
+ consumerExecutor.shutdownNow();
+ }
+
+ if (subscriber != null) {
+ subscriber.close();
+ }
+ }
+
+ @Override
+ public List poll(long durationInMillis) {
+ if (!subscribed && partitions == null) {
+ try {
+ Thread.sleep(durationInMillis);
+ } catch (InterruptedException e) {
+ log.debug("Failed to await subscription", e);
+ }
+ } else {
+ if (!subscribed) {
+ subscriptionNames = partitions.stream().map(TopicPartitionInfo::getFullTopicName).collect(Collectors.toSet());
+ subscriptionNames.forEach(admin::createTopicIfNotExists);
+ consumerExecutor = Executors.newFixedThreadPool(subscriptionNames.size());
+ messagesPerTopic = pubSubSettings.getMaxMessages()/subscriptionNames.size();
+ subscribed = true;
+ }
+ List messages;
+ try {
+ messages = receiveMessages();
+ if (!messages.isEmpty()) {
+ List result = new ArrayList<>();
+ messages.forEach(msg -> {
+ try {
+ result.add(decode(msg.getMessage()));
+ } catch (InvalidProtocolBufferException e) {
+ log.error("Failed decode record: [{}]", msg);
+ }
+ });
+ return result;
+ }
+ } catch (ExecutionException | InterruptedException e) {
+ if (stopped) {
+ log.info("[{}] Pub/Sub consumer is stopped.", topic);
+ } else {
+ log.error("Failed to receive messages", e);
+ }
+ }
+ }
+ return Collections.emptyList();
+ }
+
+ @Override
+ public void commit() {
+ acknowledgeRequests.forEach(subscriber.acknowledgeCallable()::futureCall);
+ acknowledgeRequests.clear();
+ }
+
+ private List receiveMessages() throws ExecutionException, InterruptedException {
+ List>> result = subscriptionNames.stream().map(subscriptionId -> {
+ String subscriptionName = ProjectSubscriptionName.format(pubSubSettings.getProjectId(), subscriptionId);
+ PullRequest pullRequest =
+ PullRequest.newBuilder()
+ .setMaxMessages(messagesPerTopic)
+ .setReturnImmediately(false) // return immediately if messages are not available
+ .setSubscription(subscriptionName)
+ .build();
+
+ ApiFuture pullResponseApiFuture = subscriber.pullCallable().futureCall(pullRequest);
+
+ return ApiFutures.transform(pullResponseApiFuture, pullResponse -> {
+ if (pullResponse != null && !pullResponse.getReceivedMessagesList().isEmpty()) {
+ List ackIds = new ArrayList<>();
+ for (ReceivedMessage message : pullResponse.getReceivedMessagesList()) {
+ ackIds.add(message.getAckId());
+ }
+ AcknowledgeRequest acknowledgeRequest =
+ AcknowledgeRequest.newBuilder()
+ .setSubscription(subscriptionName)
+ .addAllAckIds(ackIds)
+ .build();
+
+ acknowledgeRequests.add(acknowledgeRequest);
+ return pullResponse.getReceivedMessagesList();
+ }
+ return null;
+ }, consumerExecutor);
+
+ }).collect(Collectors.toList());
+
+ ApiFuture> transform = ApiFutures.transform(ApiFutures.allAsList(result), listMessages -> {
+ if (!CollectionUtils.isEmpty(listMessages)) {
+ return listMessages.stream().filter(Objects::nonNull).flatMap(List::stream).collect(Collectors.toList());
+ }
+ return Collections.emptyList();
+ }, consumerExecutor);
+
+ return transform.get();
+ }
+
+ public T decode(PubsubMessage message) throws InvalidProtocolBufferException {
+ DefaultTbQueueMsg msg = gson.fromJson(message.getData().toStringUtf8(), DefaultTbQueueMsg.class);
+ TbQueueMsgHeaders headers = new DefaultTbQueueMsgHeaders();
+ Map headerMap = gson.fromJson(message.getAttributesMap().get("headers"), new TypeToken
- com.amazonaws
- aws-java-sdk-sqs
- ${amazonaws.sqs.version}
-
+ com.amazonaws
+ aws-java-sdk-sqs
+ ${amazonaws.sqs.version}
+
+
+ com.google.cloud
+ google-cloud-pubsub
+ ${pubsub.client.version}
+
org.passay
passay
diff --git a/rule-engine/rule-engine-components/pom.xml b/rule-engine/rule-engine-components/pom.xml
index bf629a6f10..645419ca41 100644
--- a/rule-engine/rule-engine-components/pom.xml
+++ b/rule-engine/rule-engine-components/pom.xml
@@ -36,7 +36,6 @@
UTF-8
${basedir}/../..
1.11.747
- 1.83.0
1.16.0
@@ -99,7 +98,6 @@
com.google.cloud
google-cloud-pubsub
- ${pubsub.client.version}
com.google.api.grpc
diff --git a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml
index 6ac9b08777..a875d56782 100644
--- a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml
+++ b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml
@@ -63,7 +63,7 @@ transport:
max_string_value_length: "${JSON_MAX_STRING_VALUE_LENGTH:0}"
queue:
- type: "${TB_QUEUE_TYPE:kafka}" # kafka or aws-sqs
+ type: "${TB_QUEUE_TYPE:kafka}" # kafka or aws-sqs or pubsub
kafka:
bootstrap.servers: "${TB_KAFKA_SERVERS:localhost:9092}"
acks: "${TB_KAFKA_ACKS:all}"
@@ -75,6 +75,14 @@ queue:
access_key_id: "${TB_QUEUE_AWS_SQS_ACCESS_KEY_ID:YOUR_KEY}"
secret_access_key: "${TB_QUEUE_AWS_SQS_SECRET_ACCESS_KEY:YOUR_SECRET}"
region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}"
+ threads_per_topic: "${TB_QUEUE_AWS_SQS_THREADS_PER_TOPIC:1}"
+ visibility_timeout: "${TB_QUEUE_AWS_SQS_VISIBILITY_TIMEOUT:30}" #In seconds. If messages wont commit in this time, messages will poll again
+ pubsub:
+ project_id: "${TB_QUEUE_PUBSUB_PROJECT_ID:YOUR_PROJECT_ID}"
+ service_account: "${TB_QUEUE_PUBSUB_SERVICE_ACCOUNT:YOUR_SERVICE_ACCOUNT}"
+ ack_deadline: "${TB_QUEUE_PUBSUB_ACK_DEADLINE:30}" #In seconds. If messages wont commit in this time, messages will poll again
+ max_msg_size: "${TB_QUEUE_PUBSUB_MAX_MSG_SIZE:1048576}" #in bytes
+ max_messages: "${TB_QUEUE_PUBSUB_MAX_MESSAGES:1000}"
partitions:
hash_function_name: "${TB_QUEUE_PARTITIONS_HASH_FUNCTION_NAME:murmur3_128}"
virtual_nodes_size: "${TB_QUEUE_PARTITIONS_VIRTUAL_NODES_SIZE:16}"