|
|
|
@ -236,16 +236,18 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan |
|
|
|
edgeEvents.add(edgeEvent); |
|
|
|
} |
|
|
|
List<DownlinkMsg> 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() { |
|
|
|
|