liusuyi
2026-05-30 f4f4fc53260eb67483dce406a963628273786a61
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
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);
        }
    }
 
}