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 9c676d2e6e..b7127d4ee2 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 @@ -84,6 +84,7 @@ import java.io.Closeable; import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.Objects; import java.util.Optional; import java.util.UUID; import java.util.concurrent.CountDownLatch; @@ -393,6 +394,7 @@ public final class EdgeGrpcSession implements Closeable { return edgeEvents .stream() .map(this::convertToDownlinkMsg) + .filter(Objects::nonNull) .collect(Collectors.toList()); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryEdgeProcessor.java index dcced3d687..419cd95179 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryEdgeProcessor.java @@ -288,7 +288,7 @@ public class TelemetryEdgeProcessor extends BaseEdgeProcessor { } public DownlinkMsg processTelemetryMessageToEdge(EdgeEvent edgeEvent) throws JsonProcessingException { - EntityId entityId = null; + EntityId entityId; switch (edgeEvent.getType()) { case DEVICE: entityId = new DeviceId(edgeEvent.getEntityId()); @@ -311,12 +311,11 @@ public class TelemetryEdgeProcessor extends BaseEdgeProcessor { case EDGE: entityId = new EdgeId(edgeEvent.getEntityId()); break; + default: + log.warn("Unsupported edge event type [{}]", edgeEvent); + return null; } - DownlinkMsg downlinkMsg = null; - if (entityId != null) { - downlinkMsg = constructEntityDataProtoMsg(entityId, edgeEvent.getAction(), JsonUtils.parse(mapper.writeValueAsString(edgeEvent.getBody()))); - } - return downlinkMsg; + return constructEntityDataProtoMsg(entityId, edgeEvent.getAction(), JsonUtils.parse(mapper.writeValueAsString(edgeEvent.getBody()))); } private DownlinkMsg constructEntityDataProtoMsg(EntityId entityId, EdgeEventActionType actionType, JsonElement entityData) {