package com.ard.work.sdk.fjr.netty;
|
|
|
import com.ard.work.api.domian.ArdCamera;
|
import io.netty.bootstrap.Bootstrap;
|
import io.netty.channel.*;
|
import io.netty.channel.nio.NioEventLoopGroup;
|
import io.netty.channel.socket.nio.NioDatagramChannel;
|
import jakarta.annotation.PreDestroy;
|
import lombok.extern.slf4j.Slf4j;
|
import org.springframework.stereotype.Component;
|
import java.net.InetSocketAddress;
|
import java.util.concurrent.ConcurrentHashMap;
|
import java.util.concurrent.TimeUnit;
|
|
/**
|
* Netty 客户端管理器 - Spring 托管版本
|
*/
|
@Slf4j
|
@Component // 1. 交给 Spring 管理,默认单例
|
public class FjrNettyClient {
|
|
// 存储 CameraID 和 Channel 的映射
|
private final ConcurrentHashMap<String, Channel> channelMap = new ConcurrentHashMap<>();
|
|
// Netty 核心组件
|
private final EventLoopGroup group;
|
private final Bootstrap bootstrap;
|
|
// 2. 构造函数注入(Spring 会自动调用)
|
public FjrNettyClient() {
|
// 初始化 EventLoopGroup
|
// 注意:UDP 客户端通常只需要少量线程,这里设为 1 或 4 均可
|
this.group = new NioEventLoopGroup(4);
|
|
// 初始化 Bootstrap 公共配置
|
this.bootstrap = new Bootstrap();
|
this.bootstrap.group(group)
|
.channel(NioDatagramChannel.class)
|
.option(ChannelOption.SO_RCVBUF, 1024 * 1024) // 优化 UDP 接收缓冲
|
.handler(new ChannelInitializer<NioDatagramChannel>() {
|
@Override
|
protected void initChannel(NioDatagramChannel ch) {
|
// 添加协议编码器
|
ch.pipeline().addLast(new FjrProtocolEncoder());
|
// 可以在这里添加解码器处理服务器返回的消息
|
}
|
});
|
|
System.out.println("🚀 Netty Client Manager 初始化完成...");
|
}
|
|
/**
|
* 连接设备
|
*/
|
public boolean connect(ArdCamera camera) {
|
String cameraId = camera.getId();
|
String ip = camera.getIp();
|
int port = camera.getPort();
|
|
if (channelMap.containsKey(cameraId)) {
|
System.out.println("⚠️ 相机 " + cameraId + " 已连接,跳过");
|
return true;
|
}
|
|
try {
|
InetSocketAddress remoteAddr = new InetSocketAddress(ip, port);
|
|
// 创建新的 Channel
|
ChannelFuture future = bootstrap.bind(0).sync();
|
// 绑定远程地址(UDP Connect 只是绑定默认发送目标)
|
future.channel().connect(remoteAddr).sync();
|
|
// 存入 Map
|
channelMap.put(cameraId, future.channel());
|
|
log.info("✅ 【{}】连接成功: {}:{}", cameraId, ip, port);
|
return true;
|
} catch (Exception e) {
|
e.printStackTrace();
|
return false;
|
}
|
}
|
|
/**
|
* 发送指令
|
*/
|
public void send(String cameraId, Object msg) {
|
Channel channel = channelMap.get(cameraId);
|
if (channel != null && channel.isActive()) {
|
channel.writeAndFlush(msg);
|
} else {
|
log.error("❌ 发送失败: 相机【{}】 未连接或通道不可用", cameraId);
|
}
|
}
|
|
/**
|
* 延迟发送指令
|
*/
|
public void sendDelay(String cameraId, ThermalCommand command, long delayMs) {
|
// 使用 group 的 schedule 方法,确保在 IO 线程执行
|
group.schedule(() -> {
|
this.send(cameraId, command);
|
}, delayMs, TimeUnit.MILLISECONDS);
|
}
|
|
/**
|
* 断开指定相机
|
*/
|
public void disconnect(String cameraId) {
|
Channel channel = channelMap.remove(cameraId);
|
if (channel != null) {
|
channel.close();
|
log.info("【{}】断开连接", cameraId);
|
}
|
}
|
|
/**
|
* 3. Spring 容器关闭前的回调
|
* 确保 Netty 资源被优雅释放
|
*/
|
@PreDestroy
|
public void shutdown() {
|
System.out.println("🛑 正在关闭 Netty Client Manager...");
|
// 关闭所有连接
|
channelMap.forEach((k, v) -> v.close());
|
channelMap.clear();
|
// 关闭线程组
|
group.shutdownGracefully();
|
System.out.println("✅ Netty Client Manager 已关闭");
|
}
|
}
|