新疆aed迁移

This commit is contained in:
wanghao
2025-11-27 17:36:05 +08:00
parent fe0e075220
commit f09b695903
@@ -28,11 +28,11 @@ import java.util.concurrent.atomic.AtomicReference;
@Slf4j
public class WebSocketClientManagerNew {
@Value("${websocket.url}")
private String url;
// @Value("${websocket.url}")
// private String url;
@Value("${websocket.username}")
private String username;
// @Value("${websocket.username}")
// private String username;
@Autowired
private AedEquipmentMangerMapper aedEquipmentMangerMapper;
@@ -45,68 +45,68 @@ public class WebSocketClientManagerNew {
// connect();
}
private void connect() {
try {
log.info("aed websocket 开始建立连接...");
WebSocketContainer container = ContainerProvider.getWebSocketContainer();
String fullUrl = url + URLEncoder.encode(username, "UTF-8");
// 创建被注解的端点实例
MyEndpoint endpoint = new MyEndpoint();
// 建立连接
container.connectToServer(endpoint, new URI(fullUrl));
isReconnecting = true;
log.info("aed websocket 连接成功...");
} catch (Exception e) {
log.error("aed websocket 连接失败...", e);
scheduleReconnect();
}
}
// private void connect() {
// try {
// log.info("aed websocket 开始建立连接...");
// WebSocketContainer container = ContainerProvider.getWebSocketContainer();
// String fullUrl = url + URLEncoder.encode(username, "UTF-8");
// // 创建被注解的端点实例
// MyEndpoint endpoint = new MyEndpoint();
// // 建立连接
// container.connectToServer(endpoint, new URI(fullUrl));
// isReconnecting = true;
// log.info("aed websocket 连接成功...");
// } catch (Exception e) {
// log.error("aed websocket 连接失败...", e);
// scheduleReconnect();
// }
// }
// @ClientEndpoint
private class MyEndpoint extends Endpoint {
@Override
public void onOpen(Session session, EndpointConfig config) {
sessionRef.set(session);
try {
session.addMessageHandler(new MessageHandler.Whole<String>() {
@Override
public void onMessage(String message) {
handleMessage(message);
}
});
} catch (Exception ignore) {
log.error("Failed to add message handler");
}
}
// private class MyEndpoint extends Endpoint {
// @Override
// public void onOpen(Session session, EndpointConfig config) {
// sessionRef.set(session);
// try {
// session.addMessageHandler(new MessageHandler.Whole<String>() {
// @Override
// public void onMessage(String message) {
// handleMessage(message);
// }
// });
// } catch (Exception ignore) {
// log.error("Failed to add message handler");
// }
// }
//
// @Override
// public void onClose(Session session, CloseReason closeReason) {
// log.info("Connection closed: {}", closeReason);
// scheduleReconnect();
// }
//
// @Override
// public void onError(Session session, Throwable thr) {
// log.error("WebSocket error", thr);
// scheduleReconnect();
// }
// }
@Override
public void onClose(Session session, CloseReason closeReason) {
log.info("Connection closed: {}", closeReason);
scheduleReconnect();
}
@Override
public void onError(Session session, Throwable thr) {
log.error("WebSocket error", thr);
scheduleReconnect();
}
}
private void scheduleReconnect() {
if (isReconnecting) {
return;
}
isReconnecting = true;
scheduler.schedule(() -> {
try {
connect();
} catch (Exception e) {
log.error("aed webSocket 连接失败...", e);
} finally {
isReconnecting = false;
}
}, 20, TimeUnit.SECONDS);
}
// private void scheduleReconnect() {
// if (isReconnecting) {
// return;
// }
// isReconnecting = true;
// scheduler.schedule(() -> {
// try {
// connect();
// } catch (Exception e) {
// log.error("aed webSocket 连接失败...", e);
// } finally {
// isReconnecting = false;
// }
// }, 20, TimeUnit.SECONDS);
// }
private void handleMessage(String message) {
if (StrUtil.isBlank(message)) {