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<String, Channel> 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<SocketChannel>() {
|
@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<String> 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);
|
}
|
}
|
|
}
|