Browse Source

fixed incorrect ts in timeseries update msgs from cloud

pull/2436/head
Bohdan Smetaniuk 6 years ago
parent
commit
224c0d87e4
  1. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java
  2. 8
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java
  3. 29
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java

@ -41,7 +41,9 @@ public class EntityDataMsgConstructor {
switch (actionType) {
case TIMESERIES_UPDATED:
try {
builder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(entityData));
JsonObject data = entityData.getAsJsonObject();
long ts = data.getAsJsonPrimitive("ts").getAsLong();
builder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data.getAsJsonObject("data"), ts));
} catch (Exception e) {
log.warn("Can't convert to telemetry proto, entityData [{}]", entityData, e);
}

8
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java

@ -68,12 +68,16 @@ public class JsonConverter {
private static int maxStringValueLength = 0;
public static PostTelemetryMsg convertToTelemetryProto(JsonElement jsonElement) throws JsonSyntaxException {
public static PostTelemetryMsg convertToTelemetryProto(JsonElement jsonElement, long ts) throws JsonSyntaxException {
PostTelemetryMsg.Builder builder = PostTelemetryMsg.newBuilder();
convertToTelemetry(jsonElement, System.currentTimeMillis(), null, builder);
convertToTelemetry(jsonElement, ts, null, builder);
return builder.build();
}
public static PostTelemetryMsg convertToTelemetryProto(JsonElement jsonElement) throws JsonSyntaxException {
return convertToTelemetryProto(jsonElement, System.currentTimeMillis());
}
private static void convertToTelemetry(JsonElement jsonElement, long systemTs, Map<Long, List<KvEntry>> result, PostTelemetryMsg.Builder builder) {
if (jsonElement.isJsonObject()) {
parseObject(systemTs, result, builder, jsonElement.getAsJsonObject());

29
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.edge;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.common.util.concurrent.FutureCallback;
@ -150,13 +151,7 @@ public class TbMsgPushToEdgeNode implements TbNode {
return null;
}
ActionType actionType = getActionTypeByMsgType(msg.getType());
JsonNode entityBody = null;
JsonNode data = json.readTree(msg.getData());
if (actionType.equals(ActionType.ATTRIBUTES_UPDATED) || actionType.equals(ActionType.ATTRIBUTES_DELETED)) {
entityBody = getAttributeEntityBody(actionType, data, msg.getMetaData().getData());
} else {
entityBody = data;
}
JsonNode entityBody = getEntityBody(actionType, msg.getData(), msg.getMetaData().getData());
return buildEdgeEvent(ctx.getTenantId(), actionType, msg.getOriginator().getId(), edgeEventTypeByEntityType, entityBody);
}
}
@ -171,19 +166,25 @@ public class TbMsgPushToEdgeNode implements TbNode {
return edgeEvent;
}
private JsonNode getAttributeEntityBody(ActionType actionType, JsonNode data, Map<String, String> metadata) throws JsonProcessingException {
Map<String, Object> entityData = new HashMap<>();
private JsonNode getEntityBody(ActionType actionType, String data, Map<String, String> metadata) throws JsonProcessingException {
Map<String, Object> entityBody = new HashMap<>();
JsonNode dataJson = json.readTree(data);
switch (actionType) {
case ATTRIBUTES_UPDATED:
entityData.put("kv", data);
entityBody.put("kv", dataJson);
entityBody.put("scope", metadata.get("scope"));
break;
case ATTRIBUTES_DELETED:
List<String> keys = json.treeToValue(data.get("attributes"), List.class);
entityData.put("keys", keys);
List<String> keys = json.treeToValue(dataJson.get("attributes"), List.class);
entityBody.put("keys", keys);
entityBody.put("scope", metadata.get("scope"));
break;
case TIMESERIES_UPDATED:
entityBody.put("data", dataJson);
entityBody.put("ts", metadata.get("ts"));
break;
}
entityData.put("scope", metadata.get("scope"));
return json.valueToTree(entityData);
return json.valueToTree(entityBody);
}
private UUID getUUIDFromMsgData(TbMsg msg) throws JsonProcessingException {

Loading…
Cancel
Save