From 87b23f4989062595509fb1820a1906fce5a8aae5 Mon Sep 17 00:00:00 2001 From: Valerii Sosliuk Date: Mon, 10 Aug 2020 23:34:52 +0300 Subject: [PATCH 1/3] Display sorted metadata keys in rule nodes --- ui/src/app/common/utils.service.js | 13 ++++++++++++- ui/src/app/event/event-row.directive.js | 11 +++++++++-- .../rulechain/script/node-script-test.service.js | 4 +++- 3 files changed, 24 insertions(+), 4 deletions(-) diff --git a/ui/src/app/common/utils.service.js b/ui/src/app/common/utils.service.js index 97a04bbcf2..d5e4852afd 100644 --- a/ui/src/app/common/utils.service.js +++ b/ui/src/app/common/utils.service.js @@ -150,7 +150,8 @@ function Utils($mdColorPalette, $rootScope, $window, $translate, $q, $timeout, t customTranslation: customTranslation, objToBase64: objToBase64, base64toObj: base64toObj, - loadImageAspect: loadImageAspect + loadImageAspect: loadImageAspect, + sortObjectKeys: sortObjectKeys } return service; @@ -605,4 +606,14 @@ function Utils($mdColorPalette, $rootScope, $window, $translate, $q, $timeout, t return deferred.promise; } + function sortObjectKeys(obj) { + var sortedObj = {}; + var keys = Object.keys(obj).sort(); + for (var i = 0; i < keys.length; i++) { + var key = keys[i]; + sortedObj[key] = obj[key]; + } + return sortedObj; + } + } diff --git a/ui/src/app/event/event-row.directive.js b/ui/src/app/event/event-row.directive.js index eec8a1aff6..ffd5cc4adf 100644 --- a/ui/src/app/event/event-row.directive.js +++ b/ui/src/app/event/event-row.directive.js @@ -25,7 +25,7 @@ import eventRowDebugRuleNodeTemplate from './event-row-debug-rulenode.tpl.html'; /* eslint-enable import/no-unresolved, import/default */ /*@ngInject*/ -export default function EventRowDirective($compile, $templateCache, $mdDialog, $document, types) { +export default function EventRowDirective($compile, $templateCache, $mdDialog, $document, types, utils) { var linker = function (scope, element, attrs) { @@ -71,11 +71,18 @@ export default function EventRowDirective($compile, $templateCache, $mdDialog, $ if (!contentType) { contentType = null; } + var sortedContent; + try { + sortedContent = angular.toJson(utils.sortObjectKeys(angular.fromJson(content))); + } + catch(err) { + sortedContent = content; + } $mdDialog.show({ controller: 'EventContentDialogController', controllerAs: 'vm', templateUrl: eventErrorDialogTemplate, - locals: {content: content, title: title, contentType: contentType, showingCallback: onShowingCallback}, + locals: {content: sortedContent, title: title, contentType: contentType, showingCallback: onShowingCallback}, parent: angular.element($document[0].body), fullscreen: true, targetEvent: $event, diff --git a/ui/src/app/rulechain/script/node-script-test.service.js b/ui/src/app/rulechain/script/node-script-test.service.js index ee6eea8378..ae1f0b0568 100644 --- a/ui/src/app/rulechain/script/node-script-test.service.js +++ b/ui/src/app/rulechain/script/node-script-test.service.js @@ -20,7 +20,7 @@ import nodeScriptTestTemplate from './node-script-test.tpl.html'; /* eslint-enable import/no-unresolved, import/default */ /*@ngInject*/ -export default function NodeScriptTest($q, $mdDialog, $document, ruleChainService) { +export default function NodeScriptTest($q, $mdDialog, $document, ruleChainService, utils) { var service = { testNodeScript: testNodeScript @@ -89,6 +89,8 @@ export default function NodeScriptTest($q, $mdDialog, $document, ruleChainServic deviceName: "Test Device", ts: new Date().getTime() + "" }; + } else { + metadata = utils.sortObjectKeys(metadata); } if (!msgType) { msgType = "POST_TELEMETRY_REQUEST"; From e318b193bd604eaa1e124d3f44f438838e8b20a6 Mon Sep 17 00:00:00 2001 From: Yevhen Bondarenko <56396344+YevhenBondarenko@users.noreply.github.com> Date: Tue, 11 Aug 2020 11:15:41 +0300 Subject: [PATCH 2/3] Develop/2.5.3 confluent cloud (#3259) * added other parameters for queue kafka * Added support Confluent Cloud * fix js executor kafka connection * refactored --- .../src/main/resources/thingsboard.yml | 15 ++--- .../server/queue/kafka/TbKafkaSettings.java | 37 ++++++++++-- docker/compose-utils.sh | 5 +- docker/docker-compose.confluent.yml | 57 +++++++++++++++++++ docker/queue-confluent.env | 18 ++++++ .../config/custom-environment-variables.yml | 6 ++ msa/js-executor/config/default.yml | 4 ++ msa/js-executor/queue/kafkaTemplate.js | 20 +++++-- .../src/main/resources/tb-coap-transport.yml | 8 +++ .../src/main/resources/tb-http-transport.yml | 8 +++ .../src/main/resources/tb-mqtt-transport.yml | 8 +++ 11 files changed, 166 insertions(+), 20 deletions(-) create mode 100644 docker/docker-compose.confluent.yml create mode 100644 docker/queue-confluent.env diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 2a58209560..b6894bb0e7 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -605,16 +605,13 @@ queue: max_poll_records: "${TB_QUEUE_KAFKA_MAX_POLL_RECORDS:8192}" max_partition_fetch_bytes: "${TB_QUEUE_KAFKA_MAX_PARTITION_FETCH_BYTES:16777216}" fetch_max_bytes: "${TB_QUEUE_KAFKA_FETCH_MAX_BYTES:134217728}" + use_confluent_cloud: "${TB_QUEUE_KAFKA_USE_CONFLUENT_CLOUD:false}" + confluent: + ssl.algorithm: "${TB_QUEUE_KAFKA_CONFLUENT_SSL_ALGORITHM:https}" + sasl.mechanism: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM:PLAIN}" + sasl.config: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_JAAS_CONFIG:org.apache.kafka.common.security.plain.PlainLoginModule required username=\"CLUSTER_API_KEY\" password=\"CLUSTER_API_SECRET\";}" + security.protocol: "${TB_QUEUE_KAFKA_CONFLUENT_SECURITY_PROTOCOL:SASL_SSL}" other: -# Properties for Confluent cloud -# - key: "ssl.endpoint.identification.algorithm" -# value: "https" -# - key: "sasl.mechanism" -# value: "PLAIN" -# - key: "sasl.jaas.config" -# value: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"CLUSTER_API_KEY\" password=\"CLUSTER_API_SECRET\";" -# - key: "security.protocol" -# value: "SASL_SSL" topic-properties: rule-engine: "${TB_QUEUE_KAFKA_RE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" core: "${TB_QUEUE_KAFKA_CORE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java index d9da969eb9..37e978c07d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java @@ -18,6 +18,7 @@ package org.thingsboard.server.queue.kafka; import lombok.Getter; import lombok.Setter; import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.clients.CommonClientConfigs; import org.apache.kafka.clients.producer.ProducerConfig; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; @@ -68,7 +69,22 @@ public class TbKafkaSettings { @Value("${queue.kafka.fetch_max_bytes:134217728}") @Getter - private int fetchMaxBytes; + private int fetchMaxBytes; + + @Value("${queue.kafka.use_confluent_cloud:false}") + private boolean useConfluent; + + @Value("${queue.kafka.confluent.ssl.algorithm}") + private String sslAlgorithm; + + @Value("${queue.kafka.confluent.sasl.mechanism}") + private String saslMechanism; + + @Value("${queue.kafka.confluent.sasl.config}") + private String saslConfig; + + @Value("${queue.kafka.confluent.security.protocol}") + private String securityProtocol; @Setter private List other; @@ -76,12 +92,21 @@ public class TbKafkaSettings { public Properties toProps() { Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, servers); - props.put(ProducerConfig.ACKS_CONFIG, acks); props.put(ProducerConfig.RETRIES_CONFIG, retries); - props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize); - props.put(ProducerConfig.LINGER_MS_CONFIG, lingerMs); - props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, bufferMemory); - if(other != null){ + + if (useConfluent) { + props.put("ssl.endpoint.identification.algorithm", sslAlgorithm); + props.put("sasl.mechanism", saslMechanism); + props.put("sasl.jaas.config", saslConfig); + props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol); + } else { + props.put(ProducerConfig.ACKS_CONFIG, acks); + props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize); + props.put(ProducerConfig.LINGER_MS_CONFIG, lingerMs); + props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, bufferMemory); + } + + if (other != null) { other.forEach(kv -> props.put(kv.getKey(), kv.getValue())); } return props; diff --git a/docker/compose-utils.sh b/docker/compose-utils.sh index ac8d4fce78..a1fe74ae3f 100755 --- a/docker/compose-utils.sh +++ b/docker/compose-utils.sh @@ -39,6 +39,9 @@ function additionalComposeQueueArgs() { kafka) ADDITIONAL_COMPOSE_QUEUE_ARGS="-f docker-compose.kafka.yml" ;; + confluent) + ADDITIONAL_COMPOSE_QUEUE_ARGS="-f docker-compose.confluent.yml" + ;; aws-sqs) ADDITIONAL_COMPOSE_QUEUE_ARGS="-f docker-compose.aws-sqs.yml" ;; @@ -52,7 +55,7 @@ function additionalComposeQueueArgs() { ADDITIONAL_COMPOSE_QUEUE_ARGS="-f docker-compose.service-bus.yml" ;; *) - echo "Unknown Queue service value specified: '${TB_QUEUE_TYPE}'. Should be either kafka or aws-sqs or pubsub or rabbitmq or service-bus." >&2 + echo "Unknown Queue service value specified: '${TB_QUEUE_TYPE}'. Should be either kafka or confluent or aws-sqs or pubsub or rabbitmq or service-bus." >&2 exit 1 esac echo $ADDITIONAL_COMPOSE_QUEUE_ARGS diff --git a/docker/docker-compose.confluent.yml b/docker/docker-compose.confluent.yml new file mode 100644 index 0000000000..1bbc3f96d6 --- /dev/null +++ b/docker/docker-compose.confluent.yml @@ -0,0 +1,57 @@ +# +# 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. +# + +version: '2.2' + +services: + tb-js-executor: + env_file: + - queue-confluent.env + tb-core1: + env_file: + - queue-confluent.env + depends_on: + - redis + tb-core2: + env_file: + - queue-confluent.env + depends_on: + - redis + tb-rule-engine1: + env_file: + - queue-confluent.env + depends_on: + - redis + tb-rule-engine2: + env_file: + - queue-confluent.env + depends_on: + - redis + tb-mqtt-transport1: + env_file: + - queue-confluent.env + tb-mqtt-transport2: + env_file: + - queue-confluent.env + tb-http-transport1: + env_file: + - queue-confluent.env + tb-http-transport2: + env_file: + - queue-confluent.env + tb-coap-transport: + env_file: + - queue-confluent.env diff --git a/docker/queue-confluent.env b/docker/queue-confluent.env new file mode 100644 index 0000000000..868a135de3 --- /dev/null +++ b/docker/queue-confluent.env @@ -0,0 +1,18 @@ +TB_QUEUE_TYPE=kafka + +TB_KAFKA_SERVERS=confluent.cloud:9092 +TB_QUEUE_KAFKA_REPLICATION_FACTOR=3 + +TB_QUEUE_KAFKA_USE_CONFLUENT_CLOUD=true +TB_QUEUE_KAFKA_CONFLUENT_SSL_ALGORITHM=https +TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM=PLAIN +TB_QUEUE_KAFKA_CONFLUENT_SASL_JAAS_CONFIG=org.apache.kafka.common.security.plain.PlainLoginModule required username="CLUSTER_API_KEY" password="CLUSTER_API_SECRET"; +TB_QUEUE_KAFKA_CONFLUENT_SECURITY_PROTOCOL=SASL_SSL +TB_QUEUE_KAFKA_CONFLUENT_USERNAME=CLUSTER_API_KEY +TB_QUEUE_KAFKA_CONFLUENT_PASSWORD=CLUSTER_API_SECRET + +TB_QUEUE_KAFKA_RE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 +TB_QUEUE_KAFKA_CORE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 +TB_QUEUE_KAFKA_TA_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 +TB_QUEUE_KAFKA_NOTIFICATIONS_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:1048576000 +TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:52428800;retention.bytes:104857600 diff --git a/msa/js-executor/config/custom-environment-variables.yml b/msa/js-executor/config/custom-environment-variables.yml index c573274801..ccd33177fd 100644 --- a/msa/js-executor/config/custom-environment-variables.yml +++ b/msa/js-executor/config/custom-environment-variables.yml @@ -26,6 +26,12 @@ kafka: servers: "TB_KAFKA_SERVERS" replication_factor: "TB_QUEUE_KAFKA_REPLICATION_FACTOR" topic_properties: "TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES" + use_confluent_cloud: "TB_QUEUE_KAFKA_USE_CONFLUENT_CLOUD" + confluent: + sasl: + mechanism: "TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM" + username: "TB_QUEUE_KAFKA_CONFLUENT_USERNAME" + password: "TB_QUEUE_KAFKA_CONFLUENT_PASSWORD" pubsub: project_id: "TB_QUEUE_PUBSUB_PROJECT_ID" diff --git a/msa/js-executor/config/default.yml b/msa/js-executor/config/default.yml index f42b74745f..fac1ac7e8a 100644 --- a/msa/js-executor/config/default.yml +++ b/msa/js-executor/config/default.yml @@ -26,6 +26,10 @@ kafka: servers: "localhost:9092" replication_factor: "1" topic_properties: "retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600" + use_confluent_cloud: false + confluent: + sasl: + mechanism: "PLAIN" pubsub: queue_properties: "ackDeadlineInSec:30;messageRetentionInSec:604800" diff --git a/msa/js-executor/queue/kafkaTemplate.js b/msa/js-executor/queue/kafkaTemplate.js index 33637cee81..0420173188 100644 --- a/msa/js-executor/queue/kafkaTemplate.js +++ b/msa/js-executor/queue/kafkaTemplate.js @@ -61,15 +61,27 @@ function KafkaProducer() { const kafkaBootstrapServers = config.get('kafka.bootstrap.servers'); const requestTopic = config.get('request_topic'); + const useConfluent = config.get('kafka.use_confluent_cloud'); logger.info('Kafka Bootstrap Servers: %s', kafkaBootstrapServers); logger.info('Kafka Requests Topic: %s', requestTopic); - kafkaClient = new Kafka({ + let kafkaConfig = { brokers: kafkaBootstrapServers.split(','), - logLevel: logLevel.INFO, - logCreator: KafkaJsWinstonLogCreator - }); + logLevel: logLevel.INFO, + logCreator: KafkaJsWinstonLogCreator + }; + + if (useConfluent) { + kafkaConfig['sasl'] = { + mechanism: config.get('kafka.confluent.sasl.mechanism'), + username: config.get('kafka.confluent.username'), + password: config.get('kafka.confluent.password') + }; + kafkaConfig['ssl'] = true; + } + + kafkaClient = new Kafka(kafkaConfig); parseTopicProperties(); diff --git a/transport/coap/src/main/resources/tb-coap-transport.yml b/transport/coap/src/main/resources/tb-coap-transport.yml index c34dba7c90..f04a7559b8 100644 --- a/transport/coap/src/main/resources/tb-coap-transport.yml +++ b/transport/coap/src/main/resources/tb-coap-transport.yml @@ -69,6 +69,13 @@ queue: linger.ms: "${TB_KAFKA_LINGER_MS:1}" buffer.memory: "${TB_BUFFER_MEMORY:33554432}" replication_factor: "${TB_QUEUE_KAFKA_REPLICATION_FACTOR:1}" + use_confluent_cloud: "${TB_QUEUE_KAFKA_USE_CONFLUENT_CLOUD:false}" + confluent: + ssl.algorithm: "${TB_QUEUE_KAFKA_CONFLUENT_SSL_ALGORITHM:https}" + sasl.mechanism: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM:PLAIN}" + sasl.config: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_JAAS_CONFIG:org.apache.kafka.common.security.plain.PlainLoginModule required username=\"CLUSTER_API_KEY\" password=\"CLUSTER_API_SECRET\";}" + security.protocol: "${TB_QUEUE_KAFKA_CONFLUENT_SECURITY_PROTOCOL:SASL_SSL}" + other: topic-properties: rule-engine: "${TB_QUEUE_KAFKA_RE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" core: "${TB_QUEUE_KAFKA_CORE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" @@ -76,6 +83,7 @@ queue: notifications: "${TB_QUEUE_KAFKA_NOTIFICATIONS_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" js-executor: "${TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600}" aws_sqs: + use_default_credential_provider_chain: "${TB_QUEUE_AWS_SQS_USE_DEFAULT_CREDENTIAL_PROVIDER_CHAIN:false}" 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}" diff --git a/transport/http/src/main/resources/tb-http-transport.yml b/transport/http/src/main/resources/tb-http-transport.yml index 377c66e711..a97a58cd85 100644 --- a/transport/http/src/main/resources/tb-http-transport.yml +++ b/transport/http/src/main/resources/tb-http-transport.yml @@ -62,6 +62,13 @@ queue: linger.ms: "${TB_KAFKA_LINGER_MS:1}" buffer.memory: "${TB_BUFFER_MEMORY:33554432}" replication_factor: "${TB_QUEUE_KAFKA_REPLICATION_FACTOR:1}" + use_confluent_cloud: "${TB_QUEUE_KAFKA_USE_CONFLUENT_CLOUD:false}" + confluent: + ssl.algorithm: "${TB_QUEUE_KAFKA_CONFLUENT_SSL_ALGORITHM:https}" + sasl.mechanism: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM:PLAIN}" + sasl.config: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_JAAS_CONFIG:org.apache.kafka.common.security.plain.PlainLoginModule required username=\"CLUSTER_API_KEY\" password=\"CLUSTER_API_SECRET\";}" + security.protocol: "${TB_QUEUE_KAFKA_CONFLUENT_SECURITY_PROTOCOL:SASL_SSL}" + other: topic-properties: rule-engine: "${TB_QUEUE_KAFKA_RE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" core: "${TB_QUEUE_KAFKA_CORE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" @@ -69,6 +76,7 @@ queue: notifications: "${TB_QUEUE_KAFKA_NOTIFICATIONS_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" js-executor: "${TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600}" aws_sqs: + use_default_credential_provider_chain: "${TB_QUEUE_AWS_SQS_USE_DEFAULT_CREDENTIAL_PROVIDER_CHAIN:false}" 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}" diff --git a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml index bcc27d0755..b3804e3531 100644 --- a/transport/mqtt/src/main/resources/tb-mqtt-transport.yml +++ b/transport/mqtt/src/main/resources/tb-mqtt-transport.yml @@ -90,6 +90,13 @@ queue: linger.ms: "${TB_KAFKA_LINGER_MS:1}" buffer.memory: "${TB_BUFFER_MEMORY:33554432}" replication_factor: "${TB_QUEUE_KAFKA_REPLICATION_FACTOR:1}" + use_confluent_cloud: "${TB_QUEUE_KAFKA_USE_CONFLUENT_CLOUD:false}" + confluent: + ssl.algorithm: "${TB_QUEUE_KAFKA_CONFLUENT_SSL_ALGORITHM:https}" + sasl.mechanism: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_MECHANISM:PLAIN}" + sasl.config: "${TB_QUEUE_KAFKA_CONFLUENT_SASL_JAAS_CONFIG:org.apache.kafka.common.security.plain.PlainLoginModule required username=\"CLUSTER_API_KEY\" password=\"CLUSTER_API_SECRET\";}" + security.protocol: "${TB_QUEUE_KAFKA_CONFLUENT_SECURITY_PROTOCOL:SASL_SSL}" + other: topic-properties: rule-engine: "${TB_QUEUE_KAFKA_RE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" core: "${TB_QUEUE_KAFKA_CORE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" @@ -97,6 +104,7 @@ queue: notifications: "${TB_QUEUE_KAFKA_NOTIFICATIONS_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" js-executor: "${TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600}" aws_sqs: + use_default_credential_provider_chain: "${TB_QUEUE_AWS_SQS_USE_DEFAULT_CREDENTIAL_PROVIDER_CHAIN:false}" 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}" From 86779276cc4f047a8cadd11b31b0e1233cc87131 Mon Sep 17 00:00:00 2001 From: Valerii Sosliuk Date: Sat, 30 May 2020 15:13:24 +0300 Subject: [PATCH 3/3] Use metadata keys in originator attribute/telemetry nodes --- .../rule/engine/api/util/TbNodeUtils.java | 11 +++++++++++ .../metadata/TbAbstractGetAttributesNode.java | 9 +++++---- .../rule/engine/metadata/TbGetTelemetryNode.java | 13 +++++++------ 3 files changed, 23 insertions(+), 10 deletions(-) diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/TbNodeUtils.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/TbNodeUtils.java index 871112c1b5..4e28afd93e 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/TbNodeUtils.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/TbNodeUtils.java @@ -17,11 +17,15 @@ package org.thingsboard.rule.engine.api.util; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; +import org.springframework.util.CollectionUtils; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.server.common.msg.TbMsgMetaData; +import java.util.Collections; +import java.util.List; import java.util.Map; +import java.util.stream.Collectors; /** * Created by ashvayka on 19.01.18. @@ -41,6 +45,13 @@ public class TbNodeUtils { } } + public static List processPatterns(List patterns, TbMsgMetaData metaData) { + if (!CollectionUtils.isEmpty(patterns)) { + return patterns.stream().map(p -> processPattern(p, metaData)).collect(Collectors.toList()); + } + return Collections.emptyList(); + } + public static String processPattern(String pattern, TbMsgMetaData metaData) { String result = new String(pattern); for (Map.Entry keyVal : metaData.values().entrySet()) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java index 9c46cbbc1d..919c516ef3 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java @@ -29,6 +29,7 @@ import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.KvEntry; @@ -91,10 +92,10 @@ public abstract class TbAbstractGetAttributesNode> failuresMap = new ConcurrentHashMap<>(); ListenableFuture> allFutures = Futures.allAsList( - putLatestTelemetry(ctx, entityId, msg, LATEST_TS, config.getLatestTsKeyNames(), failuresMap), - putAttrAsync(ctx, entityId, msg, CLIENT_SCOPE, config.getClientAttributeNames(), failuresMap, "cs_"), - putAttrAsync(ctx, entityId, msg, SHARED_SCOPE, config.getSharedAttributeNames(), failuresMap, "shared_"), - putAttrAsync(ctx, entityId, msg, SERVER_SCOPE, config.getServerAttributeNames(), failuresMap, "ss_") + putLatestTelemetry(ctx, entityId, msg, LATEST_TS, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg.getMetaData()), failuresMap), + putAttrAsync(ctx, entityId, msg, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg.getMetaData()), failuresMap, "cs_"), + putAttrAsync(ctx, entityId, msg, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg.getMetaData()), failuresMap, "shared_"), + putAttrAsync(ctx, entityId, msg, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg.getMetaData()), failuresMap, "ss_") ); withCallback(allFutures, i -> { if (!failuresMap.isEmpty()) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java index a3d84b25db..4d95e8a3d1 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java @@ -103,9 +103,10 @@ public class TbGetTelemetryNode implements TbNode { if (config.isUseMetadataIntervalPatterns()) { checkMetadataKeyPatterns(msg); } - ListenableFuture> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(msg)); + List keys = TbNodeUtils.processPatterns(tsKeyNames, msg.getMetaData()); + ListenableFuture> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(msg, keys)); DonAsynchron.withCallback(list, data -> { - process(data, msg); + process(data, msg, keys); ctx.tellSuccess(ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), msg.getData())); }, error -> ctx.tellFailure(msg, error), ctx.getDbCallbackExecutor()); } catch (Exception e) { @@ -118,8 +119,8 @@ public class TbGetTelemetryNode implements TbNode { public void destroy() { } - private List buildQueries(TbMsg msg) { - return tsKeyNames.stream() + private List buildQueries(TbMsg msg, List keys) { + return keys.stream() .map(key -> new BaseReadTsKvQuery(key, getInterval(msg).getStartTs(), getInterval(msg).getEndTs(), 1, limit, NONE, getOrderBy())) .collect(Collectors.toList()); } @@ -135,7 +136,7 @@ public class TbGetTelemetryNode implements TbNode { } } - private void process(List entries, TbMsg msg) { + private void process(List entries, TbMsg msg, List keys) { ObjectNode resultNode = mapper.createObjectNode(); if (FETCH_MODE_ALL.equals(fetchMode)) { entries.forEach(entry -> processArray(resultNode, entry)); @@ -143,7 +144,7 @@ public class TbGetTelemetryNode implements TbNode { entries.forEach(entry -> processSingle(resultNode, entry)); } - for (String key : tsKeyNames) { + for (String key : keys) { if (resultNode.has(key)) { msg.getMetaData().putValue(key, resultNode.get(key).toString()); }