Browse Source

Added lock for new events updates

pull/4571/head
Volodymyr Babak 5 years ago
parent
commit
f4d0aea599
  1. 56
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

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

@ -50,6 +50,7 @@ import java.io.File;
import java.io.IOException; import java.io.IOException;
import java.io.InputStream; import java.io.InputStream;
import java.util.Collections; import java.util.Collections;
import java.util.HashMap;
import java.util.Map; import java.util.Map;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
@ -59,6 +60,8 @@ import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
@Service @Service
@Slf4j @Slf4j
@ -67,7 +70,8 @@ import java.util.concurrent.TimeUnit;
public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService { public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase implements EdgeRpcService {
private final ConcurrentMap<EdgeId, EdgeGrpcSession> sessions = new ConcurrentHashMap<>(); private final ConcurrentMap<EdgeId, EdgeGrpcSession> sessions = new ConcurrentHashMap<>();
private final ConcurrentMap<EdgeId, Boolean> sessionNewEvents = 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<>(); private final ConcurrentMap<EdgeId, ScheduledFuture<?>> sessionEdgeEventChecks = new ConcurrentHashMap<>();
private static final ObjectMapper mapper = new ObjectMapper(); private static final ObjectMapper mapper = new ObjectMapper();
@ -176,7 +180,13 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId); log.info("[{}] Closing and removing session for edge [{}]", tenantId, edgeId);
session.close(); session.close();
sessions.remove(edgeId); sessions.remove(edgeId);
sessionNewEvents.remove(edgeId); Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
lock.lock();
try {
sessionNewEvents.remove(edgeId);
} finally {
lock.unlock();
}
cancelScheduleEdgeEventsCheck(edgeId); cancelScheduleEdgeEventsCheck(edgeId);
} }
} }
@ -184,16 +194,28 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
@Override @Override
public void onEdgeEvent(TenantId tenantId, EdgeId edgeId) { public void onEdgeEvent(TenantId tenantId, EdgeId edgeId) {
log.trace("[{}] onEdgeEvent [{}]", tenantId, edgeId.getId()); log.trace("[{}] onEdgeEvent [{}]", tenantId, edgeId.getId());
if (Boolean.FALSE.equals(sessionNewEvents.get(edgeId))) { Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
log.trace("[{}] set session new events flag to true [{}]", tenantId, edgeId.getId()); lock.lock();
sessionNewEvents.put(edgeId, true); try {
if (Boolean.FALSE.equals(sessionNewEvents.get(edgeId))) {
log.trace("[{}] set session new events flag to true [{}]", tenantId, edgeId.getId());
sessionNewEvents.put(edgeId, true);
}
} finally {
lock.unlock();
} }
} }
private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) { private void onEdgeConnect(EdgeId edgeId, EdgeGrpcSession edgeGrpcSession) {
log.info("[{}] edge [{}] connected successfully.", edgeGrpcSession.getSessionId(), edgeId); log.info("[{}] edge [{}] connected successfully.", edgeGrpcSession.getSessionId(), edgeId);
sessions.put(edgeId, edgeGrpcSession); sessions.put(edgeId, edgeGrpcSession);
sessionNewEvents.put(edgeId, true); Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
lock.lock();
try {
sessionNewEvents.put(edgeId, true);
} finally {
lock.unlock();
}
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true); save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, true);
save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, System.currentTimeMillis()); save(edgeId, DefaultDeviceStateService.LAST_CONNECT_TIME, System.currentTimeMillis());
cancelScheduleEdgeEventsCheck(edgeId); cancelScheduleEdgeEventsCheck(edgeId);
@ -217,10 +239,16 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
if (sessions.containsKey(edgeId)) { if (sessions.containsKey(edgeId)) {
ScheduledFuture<?> schedule = scheduler.schedule(() -> { ScheduledFuture<?> schedule = scheduler.schedule(() -> {
try { try {
if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) { Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
log.trace("[{}] Set session new events flag to false", edgeId.getId()); lock.lock();
sessionNewEvents.put(edgeId, false); try {
session.processEdgeEvents(); if (Boolean.TRUE.equals(sessionNewEvents.get(edgeId))) {
log.trace("[{}] Set session new events flag to false", edgeId.getId());
sessionNewEvents.put(edgeId, false);
session.processEdgeEvents();
}
} finally {
lock.unlock();
} }
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, session.getEdge().getId().getId(), e); log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, session.getEdge().getId().getId(), e);
@ -249,7 +277,13 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
private void onEdgeDisconnect(EdgeId edgeId) { private void onEdgeDisconnect(EdgeId edgeId) {
log.info("[{}] edge disconnected!", edgeId); log.info("[{}] edge disconnected!", edgeId);
sessions.remove(edgeId); sessions.remove(edgeId);
sessionNewEvents.remove(edgeId); Lock lock = sessionNewEventsLocks.computeIfAbsent(edgeId, id -> new ReentrantLock());
lock.lock();
try {
sessionNewEvents.remove(edgeId);
} finally {
lock.unlock();
}
save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false); save(edgeId, DefaultDeviceStateService.ACTIVITY_STATE, false);
save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, System.currentTimeMillis()); save(edgeId, DefaultDeviceStateService.LAST_DISCONNECT_TIME, System.currentTimeMillis());
cancelScheduleEdgeEventsCheck(edgeId); cancelScheduleEdgeEventsCheck(edgeId);

Loading…
Cancel
Save