Browse Source

Use metadata keys in originator attribute/telemetry nodes

pull/3287/head
Valerii Sosliuk 6 years ago
committed by Andrew Shvayka
parent
commit
86779276cc
  1. 11
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/TbNodeUtils.java
  2. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java
  3. 13
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java

11
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<String> processPatterns(List<String> 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<String,String> keyVal : metaData.values().entrySet()) {

9
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<C extends TbGetAttributesNodeC
}
ConcurrentHashMap<String, List<String>> failuresMap = new ConcurrentHashMap<>();
ListenableFuture<List<Void>> 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()) {

13
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<TsKvEntry>> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(msg));
List<String> keys = TbNodeUtils.processPatterns(tsKeyNames, msg.getMetaData());
ListenableFuture<List<TsKvEntry>> 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<ReadTsKvQuery> buildQueries(TbMsg msg) {
return tsKeyNames.stream()
private List<ReadTsKvQuery> buildQueries(TbMsg msg, List<String> 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<TsKvEntry> entries, TbMsg msg) {
private void process(List<TsKvEntry> entries, TbMsg msg, List<String> 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());
}

Loading…
Cancel
Save