Browse Source
Merge pull request #14616 from volodymyr-babak/edge-consumer-commit-fix
Fixed infinite loop on Edge Kafka consumer commit failure
pull/14628/head
Viacheslav Klimov
8 months ago
committed by
GitHub
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
1 changed files with
5 additions and
3 deletions
-
application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java
|
|
|
@ -82,16 +82,18 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession { |
|
|
|
edgeEvents.add(edgeEvent); |
|
|
|
} |
|
|
|
List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(edgeEvents); |
|
|
|
boolean isInterrupted = true; |
|
|
|
try { |
|
|
|
boolean isInterrupted = sendDownlinkMsgsPack(downlinkMsgsPack).get(); |
|
|
|
isInterrupted = sendDownlinkMsgsPack(downlinkMsgsPack).get(); |
|
|
|
if (isInterrupted) { |
|
|
|
log.debug("[{}][{}] Send downlink messages task was interrupted", tenantId, edge.getId()); |
|
|
|
} else { |
|
|
|
consumer.commit(); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("[{}][{}] Failed to process downlink messages", tenantId, edge.getId(), e); |
|
|
|
} |
|
|
|
if (!isInterrupted) { |
|
|
|
consumer.commit(); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
|