Browse Source

Copilot codereview changes

pull/14425/head
Volodymyr Babak 10 months ago
parent
commit
3a0d083610
  1. 64
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  2. 23
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java

64
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -69,6 +69,7 @@ import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
import java.io.IOException; import java.io.IOException;
import java.io.InputStream; import java.io.InputStream;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap; import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
@ -82,6 +83,7 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock; import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Consumer; import java.util.function.Consumer;
import java.util.function.Function;
import static org.thingsboard.server.service.state.DefaultDeviceStateService.ACTIVITY_STATE; import static org.thingsboard.server.service.state.DefaultDeviceStateService.ACTIVITY_STATE;
import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_CONNECT_TIME; import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_CONNECT_TIME;
@ -654,31 +656,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private void cleanupZombieSessions() { private void cleanupZombieSessions() {
try { try {
List<EdgeId> toRemove = new ArrayList<>(); tryToDestroyZombieSessions(getZombieSessions(sessions.values()), s -> sessions.remove(s.getEdge().getId()));
for (EdgeGrpcSession session : sessions.values()) { tryToDestroyZombieSessions(getZombieSessions(sessionsById.values()), s -> sessionsById.remove(s.getSessionId()));
if (session instanceof KafkaEdgeGrpcSession kafkaSession &&
!kafkaSession.isConnected() &&
kafkaSession.getConsumer() != null &&
kafkaSession.getConsumer().getConsumer() != null &&
!kafkaSession.getConsumer().getConsumer().isStopped()) {
toRemove.add(kafkaSession.getEdge().getId());
}
if (session instanceof KafkaEdgeGrpcSession kafkaSession) {
log.debug("[{}] kafkaSession.isConnected() = {}, kafkaSession.getConsumer().getConsumer().isStopped() = {}",
kafkaSession.getEdge().getId(),
kafkaSession.isConnected(),
kafkaSession.getConsumer() != null ? kafkaSession.getConsumer().getConsumer() != null ? kafkaSession.getConsumer().getConsumer().isStopped() : null : null);
}
}
for (EdgeId edgeId : toRemove) {
log.info("[{}] Destroying session for edge because edge is not connected", edgeId);
EdgeGrpcSession removed = sessions.get(edgeId);
if (removed instanceof KafkaEdgeGrpcSession kafkaSession) {
if (kafkaSession.destroy()) {
sessions.remove(edgeId);
}
}
}
zombieSessions.removeIf(zombie -> { zombieSessions.removeIf(zombie -> {
if (zombie.destroy()) { if (zombie.destroy()) {
log.info("[{}][{}] Successfully cleaned up zombie session [{}] for edge [{}].", log.info("[{}][{}] Successfully cleaned up zombie session [{}] for edge [{}].",
@ -695,4 +675,38 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
} }
} }
private List<EdgeGrpcSession> getZombieSessions(Collection<EdgeGrpcSession> sessions) {
List<EdgeGrpcSession> result = new ArrayList<>();
for (EdgeGrpcSession session : sessions) {
if (isKafkaSessionAndZombie(session)) {
result.add(session);
}
}
return result;
}
private void tryToDestroyZombieSessions(List<EdgeGrpcSession> sessionsToRemove, Function<EdgeGrpcSession, EdgeGrpcSession> removeFunc) {
for (EdgeGrpcSession toRemove : sessionsToRemove) {
log.info("[{}] Destroying session for edge because edge is not connected", toRemove.getEdge().getId());
if (toRemove.destroy()) {
removeFunc.apply(toRemove);
}
}
}
private boolean isKafkaSessionAndZombie(EdgeGrpcSession session) {
if (session instanceof KafkaEdgeGrpcSession kafkaSession) {
log.debug("[{}] kafkaSession.isConnected() = {}, kafkaSession.getConsumer().getConsumer().isStopped() = {}",
kafkaSession.getEdge().getId(),
kafkaSession.isConnected(),
kafkaSession.getConsumer() != null ? kafkaSession.getConsumer().getConsumer() != null ? kafkaSession.getConsumer().getConsumer().isStopped() : null : null);
return !kafkaSession.isConnected() &&
kafkaSession.getConsumer() != null &&
kafkaSession.getConsumer().getConsumer() != null &&
!kafkaSession.getConsumer().getConsumer().isStopped();
}
return false;
}
} }

23
application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java

@ -108,9 +108,10 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession {
} }
if (consumer == null || (consumer.getConsumer() != null && consumer.getConsumer().isStopped())) { if (consumer == null || (consumer.getConsumer() != null && consumer.getConsumer().isStopped())) {
try { try {
if (this.consumerExecutor != null && !this.consumerExecutor.isShutdown()) { if (consumerExecutor != null && !consumerExecutor.isShutdown()) {
try { try {
this.consumerExecutor.shutdown(); consumerExecutor.shutdown();
awaitConsumerTermination();
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}][{}] Failed to shutdown previous consumer executor", tenantId, edge.getId(), e); log.warn("[{}][{}] Failed to shutdown previous consumer executor", tenantId, edge.getId(), e);
} }
@ -154,21 +155,25 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession {
} }
consumer = null; consumer = null;
try { try {
if (consumerExecutor != null) { if (consumerExecutor != null && !consumerExecutor.isShutdown()) {
consumerExecutor.shutdown(); consumerExecutor.shutdown();
try { awaitConsumerTermination();
consumerExecutor.awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS);
} catch (InterruptedException ie) {
log.warn("[{}][{}] Interrupted while awaiting consumer executor termination", tenantId, edge.getId());
}
} }
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}][{}] Failed to stop edge event consumer", tenantId, edge.getId(), e); log.warn("[{}][{}] Failed to shutdown edge event consumer executor", tenantId, edge.getId(), e);
return false; return false;
} }
return true; return true;
} }
private void awaitConsumerTermination() {
try {
consumerExecutor.awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS);
} catch (InterruptedException ie) {
log.warn("[{}][{}] Interrupted while awaiting consumer executor termination", tenantId, edge.getId());
}
}
@Override @Override
public void cleanUp() { public void cleanUp() {
String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edge.getId()).getTopic(); String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edge.getId()).getTopic();

Loading…
Cancel
Save