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 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() { @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 已关闭"); } }