Browse Source

improvements

pull/4187/head
YevhenBondarenko 6 years ago
parent
commit
f68158c550
  1. 42
      application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java
  2. 29
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java
  3. 2
      application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketMsgEndpoint.java

42
application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java

@ -41,7 +41,6 @@ import org.thingsboard.server.service.telemetry.TelemetryWebSocketMsgEndpoint;
import org.thingsboard.server.service.telemetry.TelemetryWebSocketService;
import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef;
import javax.annotation.PreDestroy;
import javax.websocket.RemoteEndpoint;
import javax.websocket.SendHandler;
import javax.websocket.SendResult;
@ -55,11 +54,7 @@ import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
@Service
@TbCoreComponent
@ -99,8 +94,6 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
private ConcurrentMap<UserId, Set<String>> regularUserSessionsMap = new ConcurrentHashMap<>();
private ConcurrentMap<UserId, Set<String>> publicUserSessionsMap = new ConcurrentHashMap<>();
private ScheduledExecutorService pingExecutor = Executors.newSingleThreadScheduledExecutor();
@Override
public void handleTextMessage(WebSocketSession session, TextMessage message) {
try {
@ -135,8 +128,6 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
}
internalSessionMap.put(internalSessionId, new SessionMetaData(session, sessionRef, maxMsgQueuePerSession));
pingExecutor.scheduleWithFixedDelay(() -> internalSessionMap.get(internalSessionId).sendPing(), 10000, 10000, TimeUnit.MILLISECONDS);
externalSessionMap.put(externalSessionId, internalSessionId);
processInWebSocketService(sessionRef, SessionEvent.onEstablished());
log.info("[{}][{}][{}] Session is opened", sessionRef.getSecurityCtx().getTenantId(), externalSessionId, session.getId());
@ -206,7 +197,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
private volatile boolean isSending = false;
private final Queue<String> msgQueue;
private final AtomicLong lastActivityTime;
private volatile long lastActivityTime;
SessionMetaData(WebSocketSession session, TelemetryWebSocketSessionRef sessionRef, int maxMsgQueuePerSession) {
super();
@ -215,14 +206,14 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
this.asyncRemote = nativeSession.getAsyncRemote();
this.sessionRef = sessionRef;
this.msgQueue = new LinkedBlockingQueue<>(maxMsgQueuePerSession);
this.lastActivityTime = new AtomicLong(System.currentTimeMillis());
this.lastActivityTime = System.currentTimeMillis();
}
synchronized void sendPing() {
try {
if (System.currentTimeMillis() - lastActivityTime.get() >= pingTimeout) {
if (System.currentTimeMillis() - lastActivityTime >= pingTimeout) {
this.asyncRemote.sendPing(ByteBuffer.wrap(new byte[]{}));
lastActivityTime.set(System.currentTimeMillis());
lastActivityTime = System.currentTimeMillis();
}
} catch (Exception e) {
log.trace("[{}] Failed to send ping msg", session.getId(), e);
@ -238,7 +229,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
if (isSending) {
try {
msgQueue.add(msg);
lastActivityTime.set(System.currentTimeMillis());
lastActivityTime = System.currentTimeMillis();
} catch (RuntimeException e) {
if (log.isTraceEnabled()) {
log.trace("[{}][{}] Session closed due to queue error", sessionRef.getSecurityCtx().getTenantId(), session.getId(), e);
@ -321,6 +312,22 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
}
}
@Override
public void sendPing(TelemetryWebSocketSessionRef sessionRef) throws IOException {
String externalId = sessionRef.getSessionId();
String internalId = externalSessionMap.get(externalId);
if (internalId != null) {
SessionMetaData sessionMd = internalSessionMap.get(internalId);
if (sessionMd != null) {
sessionMd.sendPing();
} else {
log.warn("[{}][{}] Failed to find session by internal id", externalId, internalId);
}
} else {
log.warn("[{}] Failed to find session by external id", externalId);
}
}
@Override
public void close(TelemetryWebSocketSessionRef sessionRef, CloseStatus reason) throws IOException {
String externalId = sessionRef.getSessionId();
@ -430,11 +437,4 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
}
}
@PreDestroy
private void destroy() {
if (pingExecutor != null) {
pingExecutor.shutdownNow();
}
}
}

29
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java

@ -51,22 +51,20 @@ import org.thingsboard.server.service.security.ValidationResult;
import org.thingsboard.server.service.security.ValidationResultCode;
import org.thingsboard.server.service.security.model.UserPrincipal;
import org.thingsboard.server.service.security.permission.Operation;
import org.thingsboard.server.service.subscription.TbAttributeSubscription;
import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope;
import org.thingsboard.server.service.subscription.TbEntityDataSubscriptionService;
import org.thingsboard.server.service.subscription.TbLocalSubscriptionService;
import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope;
import org.thingsboard.server.service.subscription.TbAttributeSubscription;
import org.thingsboard.server.service.subscription.TbTimeseriesSubscription;
import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper;
import org.thingsboard.server.service.telemetry.cmd.v1.AttributesSubscriptionCmd;
import org.thingsboard.server.service.telemetry.cmd.v1.GetHistoryCmd;
import org.thingsboard.server.service.telemetry.cmd.v1.SubscriptionCmd;
import org.thingsboard.server.service.telemetry.cmd.v1.TelemetryPluginCmd;
import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper;
import org.thingsboard.server.service.telemetry.cmd.v1.TimeseriesSubscriptionCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataUnsubscribeCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.DataUpdate;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate;
import org.thingsboard.server.service.telemetry.cmd.v2.UnsubscribeCmd;
import org.thingsboard.server.service.telemetry.exception.UnauthorizedException;
@ -89,6 +87,8 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import java.util.stream.Collectors;
@ -151,14 +151,23 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
private ExecutorService executor;
private String serviceId;
private ScheduledExecutorService pingExecutor;
@PostConstruct
public void initExecutor() {
serviceId = serviceInfoProvider.getServiceId();
executor = Executors.newWorkStealingPool(50);
pingExecutor = Executors.newSingleThreadScheduledExecutor();
pingExecutor.scheduleWithFixedDelay(this::sendPing, 10000, 10000, TimeUnit.MILLISECONDS);
}
@PreDestroy
public void shutdownExecutor() {
if (pingExecutor != null) {
pingExecutor.shutdownNow();
}
if (executor != null) {
executor.shutdownNow();
}
@ -744,6 +753,16 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
}
}
private void sendPing() {
wsSessionsMap.values().forEach(md ->
executor.submit(() -> {
try {
msgEndpoint.sendPing(md.getSessionRef());
} catch (IOException e) {
log.warn("[{}] Failed to send ping: {}", md.getSessionRef().getSessionId(), e);
}
}));
}
private static Optional<Set<String>> getKeys(TelemetryPluginCmd cmd) {
if (!StringUtils.isEmpty(cmd.getKeys())) {

2
application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketMsgEndpoint.java

@ -26,5 +26,7 @@ public interface TelemetryWebSocketMsgEndpoint {
void send(TelemetryWebSocketSessionRef sessionRef, int subscriptionId, String msg) throws IOException;
void sendPing(TelemetryWebSocketSessionRef sessionRef) throws IOException;
void close(TelemetryWebSocketSessionRef sessionRef, CloseStatus withReason) throws IOException;
}

Loading…
Cancel
Save