package com.ard.work.sdk.zlxd.netty.tcp; import com.ard.work.api.domian.ArdCamera; import com.ard.work.sdk.zlxd.netty.message.request.BaseXmlRequest; import com.ard.work.sdk.zlxd.utils.XmlUtils; import io.netty.bootstrap.Bootstrap; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; import io.netty.channel.*; import io.netty.channel.nio.NioEventLoopGroup; import io.netty.channel.socket.SocketChannel; import io.netty.channel.socket.nio.NioSocketChannel; import io.netty.handler.timeout.IdleStateHandler; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Component; import java.nio.charset.StandardCharsets; import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; @Component @Slf4j public class NettyTcpClient { /** * 保存每个相机的 channel(key:相机ID,value:Netty Channel) */ public static final Map CAMERA_CHANNEL_MAP = new ConcurrentHashMap<>(); @Autowired private ApplicationEventPublisher eventPublisher; // 注入发布器 // Netty客户端事件循环组(处理网络IO) private EventLoopGroup eventLoopGroup; /** * 初始化EventLoopGroup(指定线程数,多相机场景建议10-20,避免默认CPU核心数线程过多) */ @PostConstruct public void init() { eventLoopGroup = new NioEventLoopGroup(10); log.info("Netty TCP客户端初始化完成,EventLoopGroup核心线程数:{}", 10); } /** * 销毁资源(优雅关闭EventLoopGroup,释放线程/句柄,清理连接/Future) */ @PreDestroy public void destroy() { // 关闭事件循环组 if (eventLoopGroup != null && !eventLoopGroup.isShutdown()) { eventLoopGroup.shutdownGracefully(1, 5, TimeUnit.SECONDS); log.info("Netty TCP客户端EventLoopGroup已优雅关闭"); } // 清理所有相机连接 CAMERA_CHANNEL_MAP.forEach((cameraId, channel) -> disconnectCamera(cameraId)); CAMERA_CHANNEL_MAP.clear(); // 清理所有未完成的Future,避免内存泄漏 ResponseFutureHolder.clearAll(); log.info("Netty TCP客户端资源已完全释放,连接/Future清理完成"); } /** * 初始化Bootstrap(核心:无任何HTTP编解码器,仅保留空闲检测+自定义处理器) */ public Bootstrap getBootstrap(ArdCamera camera) { Bootstrap bootstrap = new Bootstrap(); bootstrap.group(eventLoopGroup) .channel(NioSocketChannel.class) // TCP优化参数(禁用Nagle算法减少延迟,开启保活检测断连) .option(ChannelOption.TCP_NODELAY, true) .option(ChannelOption.SO_KEEPALIVE, true) // 连接超时时间5秒 .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 5000) .handler(new ChannelInitializer() { @Override protected void initChannel(SocketChannel ch) { ChannelPipeline pipeline = ch.pipeline(); // 1. 空闲检测:30秒读空闲/写空闲触发事件,用于检测设备断连(和设备心跳配合) pipeline.addLast(new IdleStateHandler(30, 30, 0, TimeUnit.SECONDS)); // 2. 自定义业务处理器(直接处理原始ByteBuf,手动解析/编码HTTP,解决所有编解码错乱问题) // 每次创建新实例,解决线程安全问题 pipeline.addLast("customHttpHandler", new CustomHttpChannelHandler(camera.getId(), eventPublisher)); // 打印最终Pipeline,确认无任何HTTP编解码器(仅2个组件) pipeline.names().forEach(handlerName -> log.debug("【相机{}】Pipeline最终处理器:{}", camera.getId(), handlerName)); } }); return bootstrap; } /** * 同步连接相机(先断旧连接,再建新建连,避免重复连接) * * @param camera 相机配置(ID/IP/Port) * @return true 连接成功,false 连接失败 */ public boolean connectCamera(ArdCamera camera) { String cameraId = camera.getId(); // 先断开已有连接,避免重复持有Channel if (CAMERA_CHANNEL_MAP.containsKey(cameraId)) { disconnectCamera(cameraId); } try { // 初始化Bootstrap并连接设备 Bootstrap bootstrap = getBootstrap(camera); ChannelFuture connectFuture = bootstrap.remoteAddress(camera.getIp(), camera.getPort()).connect().sync(); if (connectFuture.isSuccess()) { Channel channel = connectFuture.channel(); CAMERA_CHANNEL_MAP.put(cameraId, channel); log.info("【相机{}】TCP连接成功 | IP:{}:{}", cameraId, camera.getIp(), camera.getPort()); return true; } else { log.error("【相机{}】TCP连接失败 | 连接Future执行失败", cameraId); return false; } } catch (InterruptedException e) { Thread.currentThread().interrupt(); log.warn("【相机{}】TCP连接被中断 | IP:{}:{}", cameraId, camera.getIp(), camera.getPort(), e); } catch (Exception e) { log.warn("【相机{}】TCP连接异常 | IP:{}:{} | 原因:{}", cameraId, camera.getIp(), camera.getPort(), e.getMessage()); } return false; } /** * 主动断开相机TCP连接(优雅关闭,清理Channel/Future) * * @param cameraId 相机ID */ public void disconnectCamera(String cameraId) { Channel channel = CAMERA_CHANNEL_MAP.remove(cameraId); if (channel == null || !channel.isOpen()) { return; } try { // 优雅关闭:等待写缓冲区数据发送完成后再关闭通道 channel.close().addListener((ChannelFutureListener) future -> { if (future.isSuccess()) { log.info("【相机{}】TCP连接已优雅关闭", cameraId); } else { log.warn("【相机{}】TCP连接关闭失败 | 原因:{}", cameraId, future.cause().getMessage()); } }); } catch (Exception e) { log.warn("【相机{}】主动断开TCP连接异常", cameraId, e); } finally { // 清理该相机的所有未完成Future,避免内存泄漏 ResponseFutureHolder.clearByCamera(cameraId); } } /** * 发送XML请求并等待响应(核心:手动拼接HTTP字节流,无任何HTTP编解码器,完美适配设备报文) * 上层业务调用无感知,入参/返回值和原有方法完全一致 */ /** * 发送XML请求并等待响应(不抛异常版本) * 上层业务调用无感知,失败返回 null */ public String sendXmlAndWait(String cameraId, BaseXmlRequest baseXmlRequest) { // 1. 校验相机通道是否可用 Channel channel = CAMERA_CHANNEL_MAP.get(cameraId); if (channel == null || !channel.isActive()) { CAMERA_CHANNEL_MAP.remove(cameraId); log.error("【相机{}】TCP通道未连接或已断开,无法发送请求", cameraId); return null; } try { String xml = XmlUtils.toXml(baseXmlRequest); int seq = baseXmlRequest.getHeader().getSequenceNumber(); String uniqueSeqKey = cameraId + "_" + seq; byte[] xmlBytes = xml.getBytes(StandardCharsets.UTF_8); //log.debug("【相机{}】准备发送业务请求 | 唯一键:{} | CSeq:{} POST | XML体长度:{} | 请求体:{}", // cameraId, uniqueSeqKey, seq, xmlBytes.length, xml); CompletableFuture future = new CompletableFuture<>(); ResponseFutureHolder.put(uniqueSeqKey, future); // 拼接HTTP请求 StringBuilder httpHeader = new StringBuilder(); httpHeader.append("POST * HTTP/1.1\r\n"); httpHeader.append("CSeq: ").append(seq).append(" POST\r\n"); httpHeader.append("Content-Length: ").append(xmlBytes.length).append("\r\n"); httpHeader.append("content-type: application/xml; charset=UTF-8\r\n"); httpHeader.append("\r\n"); byte[] httpRequestBytes = (httpHeader.toString() + xml).getBytes(StandardCharsets.UTF_8); ByteBuf requestBuf = Unpooled.copiedBuffer(httpRequestBytes); channel.writeAndFlush(requestBuf).addListener((ChannelFutureListener) sendFuture -> { if (!sendFuture.isSuccess()) { String errorMsg = "请求发送失败:" + sendFuture.cause().getMessage(); log.error("【相机{}】{} | 唯一键:{}", cameraId, errorMsg, uniqueSeqKey); ResponseFutureHolder.completeExceptionally(uniqueSeqKey, sendFuture.cause()); } else { //log.debug("【相机{}】业务请求发送成功 | 唯一键:{} | 完整请求长度:{}", // cameraId, uniqueSeqKey, httpRequestBytes.length); } }); // 同步等待响应 return future.get(10000, TimeUnit.MILLISECONDS); } catch (TimeoutException e) { log.error("【相机{}】业务请求超时(10秒) | 异常:{}", cameraId, e.getMessage()); return null; } catch (Exception e) { log.error("【相机{}】业务请求处理异常 | 异常:{}", cameraId, e.getMessage(), e); return null; } finally { ResponseFutureHolder.remove(cameraId + "_" + baseXmlRequest.getHeader().getSequenceNumber()); //log.debug("【相机{}】清理业务请求Future", cameraId); } } }