From c77217eada0bd6f208d8f4bedf0a1f70fd524718 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 15 Dec 2025 16:17:32 +0200 Subject: [PATCH] Edge Events - added merge and filter duplicates --- .../service/edge/EdgeMsgConstructorUtils.java | 133 ++++++++++++++++++ .../service/edge/rpc/EdgeGrpcSession.java | 3 +- .../edge/EdgeMsgConstructorUtilsTest.java | 98 +++++++++++++ 3 files changed, 233 insertions(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java index 505b03a1cd..6b50b8ae15 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java @@ -21,6 +21,7 @@ import com.google.gson.Gson; import com.google.gson.JsonArray; import com.google.gson.JsonElement; import com.google.gson.JsonObject; +import com.google.gson.JsonParser; import com.google.gson.JsonPrimitive; import com.google.gson.reflect.TypeToken; import lombok.extern.slf4j.Slf4j; @@ -52,6 +53,7 @@ import org.thingsboard.server.common.data.asset.AssetProfile; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.domain.DomainInfo; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.id.AiModelId; import org.thingsboard.server.common.data.id.AssetId; @@ -76,6 +78,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.id.WidgetTypeId; import org.thingsboard.server.common.data.id.WidgetsBundleId; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.notification.rule.NotificationRule; import org.thingsboard.server.common.data.notification.targets.NotificationTarget; import org.thingsboard.server.common.data.notification.template.NotificationTemplate; @@ -88,6 +91,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.UserCredentials; import org.thingsboard.server.common.data.widget.WidgetTypeDetails; import org.thingsboard.server.common.data.widget.WidgetsBundle; +import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.gen.edge.v1.AiModelUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; @@ -127,11 +131,16 @@ import org.thingsboard.server.gen.edge.v1.WidgetTypeUpdateMsg; import org.thingsboard.server.gen.edge.v1.WidgetsBundleUpdateMsg; import org.thingsboard.server.gen.transport.TransportProtos; +import java.util.ArrayList; +import java.util.Comparator; +import java.util.HashMap; +import java.util.HashSet; import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Set; import java.util.UUID; +import java.util.stream.Collectors; @Slf4j public class EdgeMsgConstructorUtils { @@ -670,4 +679,128 @@ public class EdgeMsgConstructorUtils { .setIdLSB(aiModelId.getId().getLeastSignificantBits()).build(); } + public static List mergeAndFilterDownlinkDuplicates(List edgeEvents) { + try { + edgeEvents = removeDownlinkDuplicates(edgeEvents); + + List attrUpdateMsgs = new ArrayList<>(); + for (EdgeEvent edgeEvent : edgeEvents) { + if (EdgeEventActionType.ATTRIBUTES_UPDATED.equals(edgeEvent.getAction())) { + attrUpdateMsgs.add(new AttrUpdateMsg(edgeEvent.getEntityId(), edgeEvent.getBody())); + } + } + Map> latestTsByEntityAndKey = computeLatestTsByEntityAndKey(attrUpdateMsgs); + + List result = new ArrayList<>(); + for (EdgeEvent edgeEvent : edgeEvents) { + if (!EdgeEventActionType.ATTRIBUTES_UPDATED.equals(edgeEvent.getAction())) { + result.add(edgeEvent); + continue; + } + + Map latestByKey = latestTsByEntityAndKey.get(edgeEvent.getEntityId()); + JsonNode filteredBody = filterAttributesBody(edgeEvent.getBody(), latestByKey); + if (filteredBody == null) { + continue; + } + + result.add(createFilteredEdgeEvent(edgeEvent, filteredBody)); + } + + result.sort(Comparator.comparingLong(EdgeEvent::getSeqId)); + return result; + } catch (Exception e) { + log.warn("Can't merge downlink duplicates, edgeEvents [{}]", edgeEvents, e); + return edgeEvents; + } + } + + private static AttrsTs extractAttributes(JsonNode body) { + if (body == null) { + return new AttrsTs(0L, List.of()); + } + String bodyStr = JacksonUtil.toString(body); + var jsonObject = JsonParser.parseString(bodyStr).getAsJsonObject(); + long ts = jsonObject.get("ts").getAsLong(); + var kv = jsonObject.getAsJsonObject("kv"); + List attrs = JsonConverter.convertToAttributes( + JsonUtils.getJsonObject( + JsonConverter.convertToAttributesProto(kv).getKvList() + ), ts); + return new AttrsTs(ts, attrs); + } + + private static JsonNode filterAttributesBody(JsonNode body, Map latestByKey) { + if (body == null || latestByKey == null || latestByKey.isEmpty()) { + return null; + } + String bodyStr = JacksonUtil.toString(body); + JsonObject jsonObject = JsonParser.parseString(bodyStr).getAsJsonObject(); + long ts = jsonObject.get("ts").getAsLong(); + JsonObject kv = jsonObject.getAsJsonObject("kv"); + for (Iterator> it = kv.entrySet().iterator(); it.hasNext(); ) { + Map.Entry e = it.next(); + Long latestTs = latestByKey.get(e.getKey()); + if (latestTs == null || !latestTs.equals(ts)) { + it.remove(); + } + } + if (kv.isEmpty()) { + return null; + } + return JacksonUtil.toJsonNode(jsonObject.toString()); + } + + private static Map> computeLatestTsByEntityAndKey(List attrUpdateMsgs) { + Map> latestTsByEntityAndKey = new HashMap<>(); + for (AttrUpdateMsg attrUpdateMsg : attrUpdateMsgs) { + UUID entityId = attrUpdateMsg.entityId(); + AttrsTs attrsTs = extractAttributes(attrUpdateMsg.body()); + Map map = latestTsByEntityAndKey.computeIfAbsent(entityId, id -> new HashMap<>()); + long ts = attrsTs.ts(); + for (AttributeKvEntry attr : attrsTs.attrs()) { + map.merge(attr.getKey(), ts, Math::max); + } + } + return latestTsByEntityAndKey; + } + + private static EdgeEvent createFilteredEdgeEvent(EdgeEvent edgeEvent, JsonNode filteredBody) { + EdgeEvent filtered = new EdgeEvent(); + filtered.setSeqId(edgeEvent.getSeqId()); + filtered.setTenantId(edgeEvent.getTenantId()); + filtered.setEdgeId(edgeEvent.getEdgeId()); + filtered.setAction(edgeEvent.getAction()); + filtered.setEntityId(edgeEvent.getEntityId()); + filtered.setUid(edgeEvent.getUid()); + filtered.setType(edgeEvent.getType()); + filtered.setBody(filteredBody); + return filtered; + } + + private static List removeDownlinkDuplicates(List edgeEvents) { + Set seen = new HashSet<>(); + return edgeEvents.stream() + .filter(e -> seen.add(new EventKey( + e.getTenantId(), + e.getAction(), + e.getEntityId(), + e.getType().name(), + (e.getBody() != null ? e.getBody().toString() : "null")))) + .collect(Collectors.toList()); + } + + private record EventKey(TenantId tenantId, + EdgeEventActionType action, + UUID entityId, + String type, + String body) { + } + + private record AttrsTs(long ts, List attrs) { + } + + private record AttrUpdateMsg(UUID entityId, JsonNode body) { + } + } 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 14c1b5e30a..095c0ea1f7 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 @@ -651,7 +651,8 @@ public abstract class EdgeGrpcSession implements Closeable { protected List convertToDownlinkMsgsPack(List edgeEvents) { List result = new ArrayList<>(); - for (EdgeEvent edgeEvent : edgeEvents) { + List filtered = EdgeMsgConstructorUtils.mergeAndFilterDownlinkDuplicates(edgeEvents); + for (EdgeEvent edgeEvent : filtered) { log.trace("[{}][{}] converting edge event to downlink msg [{}]", tenantId, edge.getId(), edgeEvent); DownlinkMsg downlinkMsg = null; try { diff --git a/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java b/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java index 1ee6e31372..ebb8305871 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java @@ -15,10 +15,12 @@ */ package org.thingsboard.server.service.edge; +import com.fasterxml.jackson.databind.JsonNode; import lombok.extern.slf4j.Slf4j; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; import org.junit.jupiter.api.TestInstance; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.MethodSource; @@ -36,6 +38,10 @@ import org.thingsboard.rule.engine.rest.TbSendRestApiCallReplyNode; import org.thingsboard.rule.engine.telemetry.TbCalculatedFieldsNode; import org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode; import org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode; +import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.edge.EdgeEventActionType; +import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.gen.edge.v1.EdgeVersion; @@ -44,6 +50,7 @@ import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import java.lang.reflect.Constructor; import java.util.List; import java.util.Optional; +import java.util.UUID; import java.util.stream.Stream; import static org.thingsboard.server.service.edge.EdgeMsgConstructorUtils.EXCLUDED_NODES_BY_EDGE_VERSION; @@ -155,4 +162,95 @@ public class EdgeMsgConstructorUtilsTest { String.format("For EdgeVersion '%s', ruleNode '%s' should not be included.", edgeVersion, ruleNode.getType())); } + @Test + @DisplayName("mergeDownlinkDuplicates: latest per attribute key is retained and duplicates removed") + public void testMergeDownlinkDuplicates() { + UUID deviceId = UUID.randomUUID(); + UUID assetId = UUID.randomUUID(); + TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); + + var deviceAttrUpdate1 = createEdgeEvent(tenantId, 1, EdgeEventActionType.ATTRIBUTES_UPDATED, + deviceId, EdgeEventType.DEVICE, createAttrBody(1_000L, "{\"a\":1,\"b\":1,\"d\":1}")); + var deviceAttrUpdate2 = createEdgeEvent(tenantId, 2, EdgeEventActionType.ATTRIBUTES_UPDATED, + deviceId, EdgeEventType.DEVICE, createAttrBody(2_000L, "{\"a\":2,\"b\":2,\"c\":2}")); + var deviceAttrUpdate3 = createEdgeEvent(tenantId, 3, EdgeEventActionType.ATTRIBUTES_UPDATED, + deviceId, EdgeEventType.DEVICE, createAttrBody(3_000L, "{\"a\":3,\"d\":3}")); + + var deviceUpdate = createEdgeEvent(tenantId, 4, EdgeEventActionType.UPDATED, + deviceId, EdgeEventType.DEVICE, null); + var deviceUpdateDup = createEdgeEvent(tenantId, 5, EdgeEventActionType.UPDATED, + deviceId, EdgeEventType.DEVICE, null); + + var assetAttrUpdate1 = createEdgeEvent(tenantId, 6, EdgeEventActionType.ATTRIBUTES_UPDATED, + assetId, EdgeEventType.ASSET, createAttrBody(6_000L, "{\"a\":6,\"d\":6}")); + var assetAttrUpdate2 = createEdgeEvent(tenantId, 7, EdgeEventActionType.ATTRIBUTES_UPDATED, + assetId, EdgeEventType.ASSET, createAttrBody(7_000L, "{\"a\":7,\"b\":7,\"c\":7}")); + var assetAttrUpdate3 = createEdgeEvent(tenantId, 8, EdgeEventActionType.ATTRIBUTES_UPDATED, + assetId, EdgeEventType.ASSET, createAttrBody(8_000L, "{\"a\":8,\"d\":8}")); + + List input = List.of(deviceAttrUpdate1, deviceAttrUpdate2, deviceAttrUpdate3, + deviceUpdate, deviceUpdateDup, + assetAttrUpdate1, assetAttrUpdate2, assetAttrUpdate3); + List merged = EdgeMsgConstructorUtils.mergeAndFilterDownlinkDuplicates(input); + + Assertions.assertEquals(5, merged.size()); + + EdgeEvent deviceMergedAttrBC = merged.get(0); + Assertions.assertEquals(2, deviceMergedAttrBC.getSeqId()); + Assertions.assertEquals(deviceId, deviceMergedAttrBC.getEntityId()); + Assertions.assertEquals(2_000L, deviceMergedAttrBC.getBody().get("ts").asLong()); + Assertions.assertEquals(2, getIntValue(deviceMergedAttrBC.getBody(), "b")); + Assertions.assertEquals(2, getIntValue(deviceMergedAttrBC.getBody(), "c")); + Assertions.assertNull(getIntValue(deviceMergedAttrBC.getBody(), "a")); + + EdgeEvent deviceMergedAttrAD = merged.get(1); + Assertions.assertEquals(3, deviceMergedAttrAD.getSeqId()); + Assertions.assertEquals(deviceId, deviceMergedAttrAD.getEntityId()); + Assertions.assertEquals(3_000L, deviceMergedAttrAD.getBody().get("ts").asLong()); + Assertions.assertEquals(3, getIntValue(deviceMergedAttrAD.getBody(), "a")); + Assertions.assertEquals(3, getIntValue(deviceMergedAttrAD.getBody(), "d")); + + EdgeEvent mergedDeviceUpdate = merged.get(2); + Assertions.assertEquals(4, mergedDeviceUpdate.getSeqId()); + Assertions.assertEquals(EdgeEventActionType.UPDATED, mergedDeviceUpdate.getAction()); + + EdgeEvent assetMergedAttrBC = merged.get(3); + Assertions.assertEquals(7, assetMergedAttrBC.getSeqId()); + Assertions.assertEquals(assetId, assetMergedAttrBC.getEntityId()); + Assertions.assertEquals(7_000L, assetMergedAttrBC.getBody().get("ts").asLong()); + Assertions.assertEquals(7, getIntValue(assetMergedAttrBC.getBody(), "b")); + Assertions.assertEquals(7, getIntValue(assetMergedAttrBC.getBody(), "c")); + Assertions.assertNull(getIntValue(assetMergedAttrBC.getBody(), "a")); + + EdgeEvent assetMergedAttrAD = merged.get(4); + Assertions.assertEquals(8, assetMergedAttrAD.getSeqId()); + Assertions.assertEquals(assetId, assetMergedAttrAD.getEntityId()); + Assertions.assertEquals(8_000L, assetMergedAttrAD.getBody().get("ts").asLong()); + Assertions.assertEquals(8, getIntValue(assetMergedAttrAD.getBody(), "a")); + Assertions.assertEquals(8, getIntValue(assetMergedAttrAD.getBody(), "d")); + } + + private Integer getIntValue(JsonNode body, String key) { + return body.get("kv").get(key) != null ? body.get("kv").get(key).asInt() : null; + } + + private static JsonNode createAttrBody(long ts, String kvJson) { + return JacksonUtil.toJsonNode("{\"ts\":" + ts + ",\"kv\":" + kvJson + "}"); + } + + private static EdgeEvent createEdgeEvent(TenantId tenantId, + long seqId, + EdgeEventActionType action, + UUID entityId, + EdgeEventType type, + JsonNode body) { + EdgeEvent edgeEvent = new EdgeEvent(); + edgeEvent.setSeqId(seqId); + edgeEvent.setTenantId(tenantId); + edgeEvent.setAction(action); + edgeEvent.setEntityId(entityId); + edgeEvent.setType(type); + edgeEvent.setBody(body); + return edgeEvent; + } }