Browse Source

Edge Zombie session fix - added session by session id map to handle properly connect and disconnect edge events

pull/14425/head
Volodymyr Babak 10 months ago
parent
commit
e4ad8877a3
  1. 39
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  2. 20
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java

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

@ -94,6 +94,7 @@ import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAS
public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService {
private final ConcurrentMap<EdgeId, EdgeGrpcSession> sessions = new ConcurrentHashMap<>();
private final ConcurrentMap<UUID, EdgeGrpcSession> sessionsById = new ConcurrentHashMap<>();
private final ConcurrentMap<EdgeId, Lock> sessionNewEventsLocks = new ConcurrentHashMap<>();
private final Map<EdgeId, Boolean> sessionNewEvents = new HashMap<>();
private final ConcurrentMap<EdgeId, ScheduledFuture<?>> sessionEdgeEventChecks = new ConcurrentHashMap<>();
@ -283,6 +284,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
destroySession(session);
session.cleanUp();
sessions.remove(edgeId);
sessionsById.remove(session.getSessionId());
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock();
try {
@ -332,9 +334,15 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
TenantId tenantId = edge.getTenantId();
log.info("[{}][{}] edge [{}] connected successfully.", tenantId, edgeGrpcSession.getSessionId(), edgeId);
if (sessions.containsKey(edgeId)) {
destroySession(sessions.get(edgeId));
EdgeGrpcSession existing = sessions.get(edgeId);
if (existing != null) {
log.info("[{}][{}] Replacing existing session [{}] for edge [{}]", tenantId, edgeGrpcSession.getSessionId(), existing.getSessionId(), edgeId);
destroySession(existing);
sessionsById.remove(existing.getSessionId());
}
}
sessions.put(edgeId, edgeGrpcSession);
sessionsById.put(edgeGrpcSession.getSessionId(), edgeGrpcSession);
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock();
try {
@ -492,9 +500,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private void onEdgeDisconnect(Edge edge, UUID sessionId) {
EdgeId edgeId = edge.getId();
log.info("[{}][{}] edge disconnected!", edgeId, sessionId);
EdgeGrpcSession toRemove = sessions.get(edgeId);
if (toRemove.getSessionId().equals(sessionId)) {
toRemove = sessions.remove(edgeId);
EdgeGrpcSession current = sessions.get(edgeId);
if (current != null && current.getSessionId().equals(sessionId)) {
EdgeGrpcSession toRemove = sessions.remove(edgeId);
final Lock newEventLock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
newEventLock.lock();
try {
@ -503,6 +511,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
newEventLock.unlock();
}
destroySession(toRemove);
sessionsById.remove(sessionId);
TenantId tenantId = toRemove.getEdge().getTenantId();
save(tenantId, edgeId, ACTIVITY_STATE, false);
long lastDisconnectTs = System.currentTimeMillis();
@ -510,7 +519,18 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
pushRuleEngineMessage(toRemove.getEdge().getTenantId(), edge, lastDisconnectTs, TbMsgType.DISCONNECT_EVENT);
cancelScheduleEdgeEventsCheck(edgeId);
} else {
log.debug("[{}] edge session [{}] is not available anymore, nothing to remove. most probably this session is already outdated!", edgeId, sessionId);
log.info("[{}] edge session [{}] is not current anymore. Attempting to destroy it by sessionId.", edgeId, sessionId);
EdgeGrpcSession stale = sessionsById.remove(sessionId);
if (stale != null) {
try {
destroySession(stale);
log.info("[{}][{}] Successfully destroyed stale session for edge [{}]", stale.getTenantId(), sessionId, edgeId);
} catch (Exception e) {
log.warn("[{}][{}] Failed to destroy stale session for edge [{}]", stale.getTenantId(), sessionId, edgeId, e);
}
} else {
log.debug("[{}] No session found by sessionId [{}] to destroy", edgeId, sessionId);
}
}
edgeIdServiceIdCache.evict(edgeId);
}
@ -522,6 +542,9 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
session.getTenantId(), session.getEdge().getId(), session.getEdge().getName(), session.getSessionId());
zombieSessions.add(session);
}
} catch (Exception e) {
log.warn("[{}][{}] Exception during session destroy for edge [{}] with session id [{}]",
session.getTenantId(), session.getEdge().getId(), session.getEdge().getName(), session.getSessionId(), e);
}
}
@ -640,6 +663,12 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
!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);

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

@ -101,8 +101,20 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession {
@Override
public ListenableFuture<Boolean> processEdgeEvents() {
if (!isConnected() || isSyncInProgress() || isHighPriorityProcessing) {
log.warn("[{}][{}] Session is not ready (connected={}, syncInProgress={}, highPriority={}), skip starting edge event consumer",
tenantId, edge != null ? edge.getId() : null, isConnected(), isSyncInProgress(), isHighPriorityProcessing);
return Futures.immediateFuture(Boolean.FALSE);
}
if (consumer == null || (consumer.getConsumer() != null && consumer.getConsumer().isStopped())) {
try {
if (this.consumerExecutor != null && !this.consumerExecutor.isShutdown()) {
try {
this.consumerExecutor.shutdown();
} catch (Exception e) {
log.warn("[{}][{}] Failed to shutdown previous consumer executor", tenantId, edge.getId(), e);
}
}
this.consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("edge-event-consumer"));
this.consumer = QueueConsumerManager.<TbProtoQueueMsg<ToEdgeEventNotificationMsg>>builder()
.name("TB Edge events [" + edge.getId() + "]")
@ -133,6 +145,7 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession {
public boolean destroy() {
try {
if (consumer != null) {
log.info("[{}][{}] Stopping edge event consumer...", tenantId, edge != null ? edge.getId() : null);
consumer.stop();
}
} catch (Exception e) {
@ -143,9 +156,14 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession {
try {
if (consumerExecutor != null) {
consumerExecutor.shutdown();
try {
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) {
log.warn("[{}][{}] Failed to shutdown consumer executor", tenantId, edge.getId(), e);
log.warn("[{}][{}] Failed to stop edge event consumer", tenantId, edge.getId(), e);
return false;
}
return true;

Loading…
Cancel
Save