diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java index 95b9b821ec..cf354b8877 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java @@ -236,16 +236,18 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan edgeEvents.add(edgeEvent); } List downlinkMsgsPack = downlinkMessageMapper.convertToDownlinkMsgsPack(state, edgeEvents); + boolean isInterrupted = true; try { - boolean isInterrupted = session.sendDownlinkMsgsPack(downlinkMsgsPack).get(); + isInterrupted = session.sendDownlinkMsgsPack(downlinkMsgsPack).get(); if (isInterrupted) { log.debug("[{}][{}] Send downlink messages task was interrupted", tenantId, edgeId); - } else { - consumer.commit(); } } catch (Exception e) { log.error("[{}][{}] Failed to process downlink messages", tenantId, edgeId, e); } + if (!isInterrupted) { + consumer.commit(); + } } private void cancelMigrationAndProcessingInit() {