From c3aee9c91cb420126e64c5c1d3a07350aff9a986 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 31 Mar 2020 09:58:07 +0300 Subject: [PATCH] Edge queue configuration moved to config file --- .../rule_chains/edge_root_rule_chain.json | 4 +--- .../service/edge/EdgeContextComponent.java | 5 +++++ .../edge/rpc/EdgeEventStorageSettings.java | 17 +++++++++++++++ .../service/edge/rpc/EdgeGrpcService.java | 2 +- .../service/edge/rpc/EdgeGrpcSession.java | 21 ++++++++++++------- .../src/main/resources/thingsboard.yml | 6 +++++- 6 files changed, 42 insertions(+), 13 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeEventStorageSettings.java diff --git a/application/src/main/data/json/tenant/rule_chains/edge_root_rule_chain.json b/application/src/main/data/json/tenant/rule_chains/edge_root_rule_chain.json index 4f86a073d8..0d777f029f 100644 --- a/application/src/main/data/json/tenant/rule_chains/edge_root_rule_chain.json +++ b/application/src/main/data/json/tenant/rule_chains/edge_root_rule_chain.json @@ -6,9 +6,7 @@ "firstRuleNodeId": null, "root": true, "debugMode": false, - "configuration": null, - "assignedEdges": [], - "assignedEdgesIds": [] + "configuration": null }, "metadata": { "firstNodeIndex": 2, diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index 5252fb11a0..d7ce0b4965 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -29,6 +29,7 @@ import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.relation.RelationService; +import org.thingsboard.server.service.edge.rpc.EdgeEventStorageSettings; import org.thingsboard.server.service.edge.rpc.constructor.AlarmUpdateMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.AssetUpdateMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.DashboardUpdateMsgConstructor; @@ -108,4 +109,8 @@ public class EdgeContextComponent { @Lazy @Autowired private DashboardUpdateMsgConstructor dashboardUpdateMsgConstructor; + + @Lazy + @Autowired + private EdgeEventStorageSettings edgeEventStorageSettings; } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeEventStorageSettings.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeEventStorageSettings.java new file mode 100644 index 0000000000..85b0cabc85 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeEventStorageSettings.java @@ -0,0 +1,17 @@ +package org.thingsboard.server.service.edge.rpc; + + +import lombok.Data; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; + +@Component +@Data +public class EdgeEventStorageSettings { + @Value("${edges.rpc.storage.max_read_records_count}") + private int maxReadRecordsCount; + @Value("${edges.rpc.storage.no_read_records_sleep}") + private long noRecordsSleepInterval; + @Value("${edges.rpc.storage.sleep_between_batches}") + private long sleepIntervalBetweenBatches; +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java index 51cfec8227..cb90fa4345 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java @@ -56,7 +56,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase { private boolean sslEnabled; @Value("${edges.rpc.ssl.cert}") private String certFileResource; - @Value("${edges.rpc.ssl.privateKey}") + @Value("${edges.rpc.ssl.private_key}") private String privateKeyResource; @Autowired diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index de71e91ace..036db68168 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -101,6 +101,8 @@ public final class EdgeGrpcSession implements Cloneable { private static final ReentrantLock entityCreationLock = new ReentrantLock(); + private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs"; + private final UUID sessionId; private final BiConsumer sessionOpenListener; private final Consumer sessionCloseListener; @@ -163,8 +165,7 @@ public final class EdgeGrpcSession implements Cloneable { void processHandleMessages() throws ExecutionException, InterruptedException { Long queueStartTs = getQueueStartTs().get(); - // TODO: this 100 value must be changed properly - TimePageLink pageLink = new TimePageLink(30, queueStartTs + 1000, null, true); + TimePageLink pageLink = new TimePageLink(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount(), queueStartTs, null, true); TimePageData pageData; UUID ifOffset = null; do { @@ -173,10 +174,8 @@ public final class EdgeGrpcSession implements Cloneable { log.trace("[{}] [{}] event(s) are going to be processed.", this.sessionId, pageData.getData().size()); for (Event event : pageData.getData()) { log.trace("[{}] Processing event [{}]", this.sessionId, event); - EdgeQueueEntry entry; try { - entry = objectMapper.treeToValue(event.getBody(), EdgeQueueEntry.class); - + EdgeQueueEntry entry = objectMapper.treeToValue(event.getBody(), EdgeQueueEntry.class); UpdateMsgType msgType = getResponseMsgType(entry.getType()); switch (msgType) { case ENTITY_DELETED_RPC_MESSAGE: @@ -201,6 +200,11 @@ public final class EdgeGrpcSession implements Cloneable { } if (pageData.hasNext()) { pageLink = pageData.getNextPageLink(); + try { + Thread.sleep(ctx.getEdgeEventStorageSettings().getSleepIntervalBetweenBatches()); + } catch (InterruptedException e) { + log.error("Error during sleep between batches", e); + } } } while (pageData.hasNext()); @@ -209,7 +213,7 @@ public final class EdgeGrpcSession implements Cloneable { updateQueueStartTs(newStartTs); } try { - Thread.sleep(1000); + Thread.sleep(ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval()); } catch (InterruptedException e) { log.error("Error during sleep", e); } @@ -339,13 +343,14 @@ public final class EdgeGrpcSession implements Cloneable { } private void updateQueueStartTs(Long newStartTs) { - List attributes = Collections.singletonList(new BaseAttributeKvEntry(new LongDataEntry("queueStartTs", newStartTs), System.currentTimeMillis())); + newStartTs = ++newStartTs; // increments ts by 1 - next edge event search starts from current offset + 1 + List attributes = Collections.singletonList(new BaseAttributeKvEntry(new LongDataEntry(QUEUE_START_TS_ATTR_KEY, newStartTs), System.currentTimeMillis())); ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes); } private ListenableFuture getQueueStartTs() { ListenableFuture> future = - ctx.getAttributesService().find(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, "queueStartTs"); + ctx.getAttributesService().find(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, QUEUE_START_TS_ATTR_KEY); return Futures.transform(future, attributeKvEntryOpt -> { if (attributeKvEntryOpt != null && attributeKvEntryOpt.isPresent()) { AttributeKvEntry attributeKvEntry = attributeKvEntryOpt.get(); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 869bdb09cd..2975076483 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -571,7 +571,11 @@ edges: # Enable/disable SSL support enabled: "${EDGES_RPC_SSL_ENABLED:false}" cert: "${EDGES_RPC_SSL_CERT:certChainFile.pem}" - privateKey: "${EDGES_RPC_SSL_PRIVATE_KEY:privateKeyFile.pem}" + private_key: "${EDGES_RPC_SSL_PRIVATE_KEY:privateKeyFile.pem}" + storage: + max_read_records_count: "${EDGES_RPC_STORAGE_MAX_READ_RECORDS_COUNT:50}" + no_read_records_sleep: "${EDGES_RPC_NO_READ_RECORDS_SLEEP:1000}" + sleep_between_batches: "${EDGES_RPC_SLEEP_BETWEEN_BATCHES:1000}" swagger: api_path_regex: "${SWAGGER_API_PATH_REGEX:/api.*}"