From 70ba4b03e4f1d8cb81a25c251027a560f15bea6b Mon Sep 17 00:00:00 2001 From: Nikita Mazurenko Date: Tue, 16 Dec 2025 19:01:12 +0200 Subject: [PATCH] Fix infinite loop on Kafka consumer commit failure --- .../session/manager/KafkaBasedEdgeGrpcSessionManager.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) 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() {