--- title: 网络编程简介 date: 2025-12-15 --- # 一、TCP ## 1.1 定义   TCP(Transmission Control Protocol,传输控制协议)是互联网核心的**面向连接、可靠、字节流**的传输层协议,工作在 OSI 模型的传输层(TCP/IP 模型的传输层),基于 IP 协议提供端到端的可靠数据传输服务。它是 HTTP、HTTPS、WebSocket、MQTT 等应用层协议的底层依赖,核心目标是解决 IP 协议 “无连接、不可靠、无顺序” 的缺陷,确保数据在不可靠的网络中准确、完整、有序地传输。 ## 1.2 特性 ### 1.2.1 面向连接   通信前必须完成「三次握手」建立连接,通信后通过「四次挥手」释放连接: - 三次握手:客户端发 SYN → 服务器回 SYN+ACK → 客户端发 ACK(确保双方收发能力正常); - 四次挥手:客户端发 FIN → 服务器回 ACK → 服务器发 FIN → 客户端回 ACK(确保数据传输完毕)。 ### 1.2.2 可靠传输 - 序号与确认号:每个字节都有序号,接收方收到后回复确认号,未收到则发送方重传; - 超时重传:发送方未在规定时间收到确认,自动重传数据; - 流量控制:通过滑动窗口机制,防止发送方发送过快导致接收方缓冲区溢出; - 拥塞控制:通过慢启动、拥塞避免等算法,适应网络带宽变化。 ### 1.2.3 面向字节流   TCP 将应用层数据视为连续的字节流,不保留应用层数据的边界(与 UDP 的 “数据报” 模式不同): - 发送方:应用层数据被拆分为 TCP 报文段(Segment)发送,拆分规则由 TCP 协议决定(如 MSS 限制)。 - 接收方:将收到的报文段按顺序重组为完整的字节流,再交给应用层,确保数据顺序与发送时一致。 ### 1.2.4 有序传输   TCP 报文段头部包含 “序号(Sequence Number)” 和 “确认号(Acknowledgment Number)”: - 序号(SN):标识发送方当前发送的字节流位置(如序号为 100 表示当前报文段的第一个字节是整个字节流的第 100 字节)。 - 确认号(ACK):标识接收方期望下次接收的字节流位置(如确认号为 200 表示已正确接收前 199 字节,下次需从 200 字节开始接收)。 - 接收方通过序号排序报文段,丢弃重复报文,确保按发送顺序交付数据。 ### 1.2.5 全双工通信   TCP 连接是双向的,双方可同时发送和接收数据,无需等待对方结束发送: - 每个方向都有独立的发送缓冲区和接收缓冲区,以及独立的滑动窗口用于流量控制。 - 示例:客户端发送数据的同时,服务器可同步向客户端返回响应,无需等待客户端发送完毕。 ## 1.3 优缺点   优点: 1. 可靠、有序、无丢包 2. 支持流量 / 拥塞控制 3. 适用于大数据传输   缺点: 1. 连接建立 / 释放开销大 2. 头部开销大(20-60 字节) 3. 不适合实时性要求极高的场景(如直播低延迟) ## 1.4 Python实现   TCP服务器: ```python import socket import signal import threading from concurrent.futures import ThreadPoolExecutor from dotenv import load_dotenv import os from model.mqtt import UNKNOWN_MESSAGE from processor.tcp_processor import TcpMessageProcessor from config.logger_config import logger class TCPServer: def __init__(self): load_dotenv() self.host = os.getenv("TCP_HOST", 'localhost') self.port = int(os.getenv("TCP_PORT", 9100)) self.max_workers = 10 self.timeout = 300 self.server_socket = None self.running = False self.thread_pool = ThreadPoolExecutor(max_workers=self.max_workers) self.message_processor = TcpMessageProcessor() # 已经连接的客户端 self.clients = {} self.client_lock = threading.Lock() self.server_thread = None signal.signal(signal.SIGTERM, self._handle_signal) signal.signal(signal.SIGINT, self._handle_signal) def _handle_signal(self, signum, frame): """处理终止信号,触发优雅关闭""" logger.info(f"收到信号 {signum},准备关闭服务器...") self.running = False def _start_loop(self): """启动服务器""" try: # 创建 TCP/IP socket AF_INET: IPv4地址族 SOCK_STREAM: TCP协议(面向连接) self.server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) # 设置SO_REUSEADDR选项,允许重用地址和端口,避免服务端重启时出现 “地址已被占用” 的错误 self.server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) # 绑定到指定主机和端口 0.0.0.0 表示所有主机 self.server_socket.bind((self.host, self.port)) # 设置最大等待连接数 self.server_socket.listen(5) # 设置socket超时时间(1.0秒) self.server_socket.settimeout(1.0) self.running = True logger.info(f"TCP服务器启动,监听 {self.host}:{self.port} " f"(最大线程: {self.max_workers}, 超时: {self.timeout}s)") while self.running: try: # accept() 会阻塞直到有客户端连接 # client_socket: 与客户端通信的新socket client_address: 客户端地址(ip, port)元组 client_socket, client_address = self.server_socket.accept() # 获取客户端IP client_ip = client_address[0] # 设置客户端socket超时 # 防止客户端长时间不发送数据 client_socket.settimeout(self.timeout) logger.info(f"新连接: {client_address}") # 存储客户端 with self.client_lock: self.clients[client_ip] = client_socket # 提交到线程池处理 # 提交到线程池是为了让服务器能同时服务多个客户端,而不让一个慢客户端阻塞所有其他客户端 self.thread_pool.submit(self.handle_client, client_socket, client_ip) except socket.timeout: continue except Exception as e: if self.running: logger.error(f"接受连接失败: {str(e)}") except Exception as e: logger.error(f"TCP服务器启动失败: {str(e)}") finally: self.stop() def start(self): self.server_thread = threading.Thread(target=self._start_loop, daemon=True) self.server_thread.start() def stop(self): if not self.running: return self.running = False logger.info("开始关闭服务器...") # 移除客户端 with self.client_lock: for client_socket in self.clients.values(): try: client_socket.close() except Exception as e: logger.warning(f"关闭客户端连接失败: {str(e)}") self.clients.clear() # 关闭线程池 self.thread_pool.shutdown(wait=True) logger.info("所有客户端处理线程已结束") # 关闭连接 if self.server_socket: self.server_socket.close() logger.info(f"服务器已关闭({self.host}:{self.port})") def handle_client(self, client_socket, client_ip): """处理客户端连接""" try: while True: data = client_socket.recv(1024) if not data: logger.info(f"客户端 {client_ip} 主动断开连接") break message = data.decode('utf-8').strip() logger.info(f"收到 {client_ip} 的消息: {message}") # 放到消息处理器里面处理 response = self.message_processor.process(message) # 回复消息 if response != UNKNOWN_MESSAGE: client_socket.sendall(response.encode('utf-8')) logger.info(f"回复 {client_ip}: {response}") except socket.timeout: logger.warning(f"客户端 {client_ip} 超时未活动") except Exception as e: logger.error(f"处理 {client_ip} 出错: {str(e)}") finally: # 异常情况下关闭连接 with self.client_lock: if client_ip in self.clients: del self.clients[client_ip] try: client_socket.close() logger.info(f"客户端 {client_ip} 连接已关闭") except Exception as e: logger.warning(f"关闭 {client_ip} 连接失败: {str(e)}") def send_to_client(self, client_ip, message): # 先获取客户端连接(加锁保护) with self.client_lock: client_socket = self.clients.get(client_ip) if not client_socket: logger.warning(f"客户端 {client_ip} 不存在或已断开连接") return False # 发送消息 try: client_socket.sendall(message.encode('utf-8')) logger.info(f"主动发送消息给 {client_ip}: {message}") return True except Exception as e: logger.error(f"向 {client_ip} 发送消息失败: {str(e)}") # 发送失败时移除无效连接 with self.client_lock: if client_ip in self.clients: del self.clients[client_ip] return False tcp_server = TCPServer() ```   main.py: ```python from endpoint.tcp_server import tcp_server from config.logger_config import logger import time if __name__ == "__main__": try: tcp_server.start() try: while True: time.sleep(1) except KeyboardInterrupt: logger.info("收到终止信号,开始关闭程序...") finally: tcp_server.stop() logger.info("程序已退出") except Exception as e: logger.critical(f"程序启动失败: {str(e)}", exc_info=True) exit(1) ```   启动tcp_server后,通过while循环防止主线程退出,从而让后台的TCP服务器线程能继续运行。   在TCPServer中,通过while循环持续接收tcp客户端的连接,每当有一个客户端连接时,会提交到线程池中去处理,通过自定义消息处理器,将处理完的结果返回给客户端。   TCP消息处理器: ::: code-group ```python [抽象消息处理器] from abc import ABC, abstractmethod class AbstractDeviceMessageParser(ABC): """设备消息解析器基类""" @abstractmethod def check(self, message): """判断当前解析器是否能处理该消息""" pass @abstractmethod def parse(self, message): """解析消息内容""" pass @abstractmethod def response(self, parsed_data): """根据解析后的数据生成响应""" pass ``` ```python [示例消息处理器] from datetime import datetime from model.band import BandData from model.mqtt import MqttTopic, MqttData, UNKNOWN_MESSAGE from parser.abstract_parsers import AbstractDeviceMessageParser import re from utils.index import str_length_to_4hex from config.logger_config import logger from endpoint.mqtt_client import mqtt_client class JuweiBandParser(AbstractDeviceMessageParser): """聚伟手环消息解析器""" def __init__(self): # 聚伟手环消息格式: [MNYD*设备ID*内容长度*内容] self.pattern = r'MNYD' self.vendor = "聚伟手环" self.tag = "MNYD" self.device_id = "" def check(self, message): """检查是否为聚伟手环的消息格式""" return re.search(self.pattern, message) is not None def parse(self, message): """解析聚伟手环消息""" try: parts = message.strip("[]").split("*") self.device_id = parts[1] content = parts[3] return { 'vendor': self.vendor, 'device_id': self.device_id, 'content': content, } except Exception as e: logger.error(f"解析{self.vendor}消息出错: {str(e)}") return None def publish_vital_data(self, item, value): mqtt_client.publish( topic=MqttTopic.JuWei_Band_Post.value, payload=MqttData( deviceIp="", deviceId=self.device_id, payload=BandData(item=item, value=value).model_dump_json(), ).model_dump_json()) def response(self, parsed_data): """生成聚伟手环的响应消息""" if not parsed_data: return UNKNOWN_MESSAGE # 处理不同命令 content = parsed_data['content'] parts = content.split(",", 1) content_tag = parts[0] if len(parts) > 0 else "" content_value = parts[1] if len(parts) > 1 else "" replay = UNKNOWN_MESSAGE logger.info(f"解析{content_tag}消息") match content_tag: # PING消息 [MNYD*334588000000156*0004*PING] case "PING": replay = "PING,1" # 日期,步数,翻滚次数,电量百分数,里程数(km) # [MNYD*334588000000156*0014*KA,120414,50,100,100,100.12] case "KA": replay = content_tag # 位置数据上报 # case "UD": # replay = "" # 报警数据上报 # [MNYD*334588000000156*00CD*AL,180916,064153,A,22.570512,N,113.8623267, # E,0.00,154.8,0.0,11,100,100,0,0,00100018,7,0,460,1,9529, # 21809,155,9529,21242,132,9529,21405,131,9529,63554,131,9529, # 63555,130,9529,63556,118,9529,21869,116,0,12.4] case "AL": replay = content_tag # 获取服务器端时间 # [MNYD*YYYYYYYYYYYYYYY*LEN*LGZONE] case "LGZONE": now = datetime.now() current_date = now.date().strftime("%Y-%m-%d") current_time = now.time().strftime("%H:%M:%S") replay = f"{content_tag},+8,{current_time},{current_date}" # 请求位置数据 TODO case "WG": replay = content_tag # 请求电话本设置信息 TODO case "PHLQ": replay = "PHL" # 请求SOS设置信息 TODO case "SOS": replay = content_tag # 终端心率上传 case "heart": self.publish_vital_data(content_tag, content_value) replay = content_tag # 上传体温数据 [MNYD*334588000000156*0009*temp,36.2] case "temp": self.publish_vital_data(content_tag, content_value) replay = content_tag # 上传血压数据 [MNYD*334588000000156*000C*blood,150,70] case "blood": self.publish_vital_data(content_tag, content_value) replay = content_tag # 上传血氧数据 [MNYD*334588000000156*0009*oxygen,97] case "oxygen": self.publish_vital_data(content_tag, content_value) replay = content_tag # 上传睡眠数据报告 case "SLEEPRPT": replay = "SLEEP" case _: logger.info("不需要回复") return UNKNOWN_MESSAGE return f"[{self.tag}*{parsed_data['device_id']}*{str_length_to_4hex(replay)}*{replay}]" ``` ```python [消息处理器] from model.mqtt import UNKNOWN_MESSAGE from parser.juwei_band_parser import JuweiBandParser from config.logger_config import logger class TcpMessageProcessor: """消息处理器,负责将消息路由到正确的设备解析器""" def __init__(self): # 注册所有支持的设备解析器 self.parsers = [JuweiBandParser()] def process(self, message): """处理消息,返回响应""" # 尝试找到能处理该消息的解析器 for parser in self.parsers: if parser.check(message): parsed_data = parser.parse(message) if parsed_data: logger.info(f"处理{parsed_data['vendor']}消息: {message}") return parser.response(parsed_data) # 没有找到合适的解析器 logger.warning(f"未识别的消息格式: {message}") return UNKNOWN_MESSAGE ``` ::: # 二、HTTP ## 2.1 定义   HTTP(HyperText Transfer Protocol,超文本传输协议)是互联网的核心协议之一,用于**客户端(如浏览器、App)与服务器之间的通信**,是万维网(WWW)数据交换的基础。它定义了请求 / 响应的格式、传输规则和状态码等核心机制,支持从简单文本到复杂多媒体(图片、视频、文件)的传输,也是现代 Web 应用的底层通信标准。 ## 2.2 特性 ### 2.2.1 请求 - 响应模式   通信由客户端主动发起请求,服务器接收后处理并返回响应,不存在服务器主动向客户端推送数据的情况(HTTP/2 引入 Server Push 扩展,可主动推送关联资源)。   一次完整通信流程:**客户端建立连接 → 发送请求 → 服务器处理 → 返回响应 → 连接关闭**(HTTP/1.1 默认开启长连接 Keep-Alive)。 ### 2.2.2 无状态   **服务器不会保存客户端的会话状态**(如登录状态、浏览记录),每次请求都是独立的,服务器无法通过协议本身识别连续请求是否来自同一客户端。通过 Cookie、Session、Token(如 JWT)等机制补充状态管理。 ## 2.3 版本 ### 2.3.1 HTTP/1.0(1996 年) - 基础版本,支持 GET、POST、HEAD 三种请求方法。 - 每次请求都需要建立新的 TCP 连接(短连接),连接建立和关闭的开销大,性能较低。 - 不支持长连接、管线化请求,仅支持简单的文本传输。 ### 2.3.2 HTTP/1.1(1999 年) - 默认开启 长连接(Keep-Alive):同一 TCP 连接可处理多个请求,减少连接开销。 - 支持 管线化请求:客户端可连续发送多个请求,无需等待前一个响应返回(部分浏览器未完全支持)。 - 新增请求方法:PUT、DELETE、OPTIONS、TRACE、CONNECT。 - 支持 chunked 编码(分块传输)、缓存控制(Cache-Control)、内容协商等核心功能。 ### 2.3.3 HTTP/2(2015 年)   基于 SPDY 协议优化,核心目标是提升性能: - 二进制帧传输:将请求 / 响应数据拆分为二进制帧,而非 HTTP/1.x 的文本格式,解析效率更高。 - 多路复用:同一 TCP 连接中可并发处理多个请求(通过帧的 Stream ID 区分),解决 HTTP/1.1 的 “队头阻塞” 问题。 - 服务器推送(Server Push):服务器可主动向客户端推送关联资源(如 HTML 引用的 CSS/JS),减少客户端请求次数。 - 头部压缩(HPACK):对请求头和响应头进行压缩,减少传输体积(HTTP/1.x 头部重复传输开销大)。 ### 2.3.4 HTTP/3(2022 年)   基于 SPDY 协议优化,核心目标是提升性能: - 解决 TCP 队头阻塞:UDP 无连接特性,单个流的阻塞不影响其他流。 - 更快的连接建立:QUIC 集成 TLS 1.3,减少握手次数(1-RTT 甚至 0-RTT 建立连接)。 - 更好的移动网络支持:支持连接迁移(如手机切换 WiFi/4G 时,连接不中断)。 ## 2.4 组成 ### 2.4.1 请求消息(Request)   客户端向服务器发送的请求格式,由 请求行、请求头、空行、请求体 四部分组成: ```http GET /api/courses/1 HTTP/1.1 # 请求行 Host: lms.example.com # 请求头(键值对形式) Authorization: Bearer Accept: application/json User-Agent: Mozilla/5.0 (Chrome/120.0.0.0) Content-Type: application/json {"username": "admin", "password": "123456"} # 请求体(可选,POST/PUT 等方法常用) ``` #### 请求行 - 请求方法:表示请求的操作类型(常用方法如下表)。 - 请求 URI:指定服务器上的资源路径。 - 协议版本:如 HTTP/1.1、HTTP/2。 #### 请求头   描述请求的附加信息,常用字段: - Host:目标服务器域名(如 lms.example.com),HTTP/1.1 必选字段。 - User-Agent:客户端身份标识(如浏览器版本、App 名称)。 - Accept:客户端可接收的响应数据格式(如 application/json、text/html)。 - Content-Type:请求体的数据格式(如 application/json、multipart/form-data(文件上传))。 - Authorization:身份认证信息(如 Token、Basic Auth)。 - Cookie:客户端存储的会话信息(如登录态 Cookie)。 - Cache-Control:缓存控制策略(如 no-cache 表示不使用缓存)。 #### 请求体   可选部分,仅在需要向服务器提交数据时使用(如 POST 提交表单、PUT 更新资源),数据格式由 Content-Type 指定: - 表单数据:application/x-www-form-urlencoded(如 username=admin&password=123)。 - JSON 数据:application/json(如 {"key": "value"})。 - 文件上传:multipart/form-data(如 LMS 系统的作业文件上传)。 - 纯文本:text/plain。 ### 2.4.2 响应消息(Response)   服务器向客户端返回的响应格式,由 状态行、响应头、空行、响应体 四部分组成: ```http HTTP/1.1 200 OK # 状态行 Server: Nginx Content-Type: application/json Content-Length: 128 Set-Cookie: sessionId=abc123; Path=/ {"code": 200, "message": "success", "data": {"id": 1, "name": "Vue3 实战课程"}} # 响应体 ``` #### 状态行 - 协议版本:如 HTTP/1.1。 - 状态码:表示请求处理结果。 - 状态短语:状态码的文字描述(如 OK、Not Found)。 #### HTTP状态码 | 分类 | 状态码范围 | 含义 | 常用码 | |------|------------|-----------------------|-------------------------| | 1xx | 100-199 | 信息性响应(临时响应)| 100 Continue(预检通过)| | 2xx | 200-299 | 成功响应 | 200 OK(成功)、201 Created(资源创建成功)、204 No Content(成功无响应体)| | 3xx | 300-399 | 重定向 | 301 永久重定向、302 临时重定向、304 Not Modified(缓存有效)| | 4xx | 400-499 | 客户端错误 | 400 Bad Request(请求参数错误)、401 Unauthorized(未认证)、403 Forbidden(权限不足)、404 Not Found(资源不存在)、405 Method Not Allowed(请求方法不支持)| | 5xx | 500-599 | 服务器错误 | 500 Internal Server Error(服务器内部错误)、502 Bad Gateway(网关错误)、503 Service Unavailable(服务不可用)、504 Gateway Timeout(网关超时)| #### 响应头   描述响应的附加信息,常用字段: - Server:服务器软件标识(如 Nginx、Tomcat)。 - Content-Type:响应体的数据格式(如 application/json、text/html)。 - Content-Length:响应体的字节大小。 - Set-Cookie:服务器向客户端设置 Cookie(如登录态、会话 ID)。 - Cache-Control:缓存控制策略(如 max-age=3600 表示缓存 1 小时)。 - Access-Control-Allow-Origin:跨域资源共享(CORS)配置(如 * 表示允许所有域名跨域)。 - Location:重定向目标地址(3xx 状态码必选)。 #### 响应体   服务器返回的核心数据,数据格式由 Content-Type 指定,常见格式: - JSON:前后端分离项目首选(如 {"code": 200, "data": [...]})。 - HTML:传统 Web 页面(如静态网页、JSP 页面)。 - 图片 / 视频:二进制流(如 image/jpeg、video/mp4)。 - 纯文本:text/plain。 # 三、WebSocket ## 3.1 定义   WebSocket 是一种**全双工、双向、持久化的网络通信协议**(属于应用层协议),由 HTML5 规范定义,专门解决 HTTP 协议无法实现服务器主动向客户端推送数据的问题。   它的核心特点是:**一旦客户端与服务器建立连接,双方就可以在这个连接上实时、双向地发送数据**,无需像 HTTP 那样每次通信都由客户端发起请求,非常适合实时通信场景(如聊天、直播弹幕、实时数据监控、在线协作等)。 | 特性 | HTTP | WebSocket | |---------------------|-------------------------------|-------------------------------| | 通信方向 | 单向(客户端请求→服务器响应)| 全双工(双方可同时发数据)| | 连接类型 | 短连接 / 长连接(需重复请求) | 持久连接(一次建立,持续通信) | | 数据传输效率 | 每次请求带大量头部信息,效率低 | 连接建立后仅传输数据,开销小 | | 服务器主动推送 | 不支持(HTTP/2 的 Server Push 仅能推送资源,非实时数据) | 原生支持,可主动向客户端发数据 | | 协议标识 | http:// / https:// | ws:// / wss://(加密版)| ## 3.2 特性 ### 3.2.1 握手过程 1. 客户端发送 HTTP 请求,请求头包含 Upgrade: websocket 和 Connection: Upgrade(表示要升级为 WebSocket 协议); 2. 服务器响应 101 Switching Protocols,握手成功,连接转为 WebSocket 持久连接; 3. 后续通信不再使用 HTTP 格式,而是 WebSocket 帧格式(二进制 / 文本)。   请求头: ```http GET /chat HTTP/1.1 Host: example.com Upgrade: websocket # 核心:请求升级为WebSocket协议 Connection: Upgrade # 核心:表示连接要升级 Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ== # 客户端生成的随机密钥,用于验证 Sec-WebSocket-Version: 13 # 指定WebSocket版本(主流为13) Origin: https://example.com # 跨域验证 ```   响应头: ```http HTTP/1.1 101 Switching Protocols Upgrade: websocket Connection: Upgrade Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo= # 服务器加密后的密钥,客户端验证 ``` ### 3.2.2 全双工通信   此时 HTTP 连接已升级为 WebSocket 连接,双方可以通过这个连接实时、双向地发送数据。数据传输采用帧(Frame) 格式,支持文本数据和二进制数据(如图片、视频流)。 ### 3.2.3 帧格式   WebSocket 数据以帧为单位传输,帧头包含操作码(文本帧 0x01、二进制帧 0x02、关闭帧 0x08 等)、掩码(客户端发送数据需掩码,服务器无需); ### 3.2.4 无同源限制   WebSocket 不遵循同源策略(但服务器可通过 Origin 头限制跨域); ### 3.2.5 心跳机制   通过 Ping/Pong 帧维持连接(避免网络设备断开空闲连接,如 LMS 系统需定期发送 Ping 帧,服务器回复 Pong 帧)。 # 四、MQTT ## 4.1 定义   MQTT(Message Queuing Telemetry Transport,消息队列遥测传输)是一种**轻量级、低带宽、低功耗的发布 / 订阅(Publish/Subscribe)模式物联网(IoT)通信协议**,由 IBM 于 1999 年设计,核心目标是解决受限设备(如传感器、嵌入式设备)和低带宽、不稳定网络环境下的高效数据传输问题。 ## 4.2 架构 - 发布者(Publisher):发送消息的设备 / 服务; - 订阅者(Subscriber):接收消息的设备 / 服务; - broker(代理服务器):核心中间件,接收发布者的消息,根据「主题(Topic)」转发给订阅者(如 EMQ X、Mosquitto、RabbitMQ 支持 MQTT 插件); - 主题(Topic):消息的分类标识。 ## 4.3 特性 ### 4.3.1 QoS(服务质量)等级 - QoS 0(最多一次):消息发送一次,不保证送达; - QoS 1(至少一次):消息至少送达一次,可能重复; - QoS 2(恰好一次):消息仅送达一次,最可靠。 ### 4.3.2 轻量级   头部开销极小(固定头部仅 2 字节),消息体支持二进制 / 文本,适合低带宽场景 ### 4.3.3 保留消息(Retained Message)   broker 保存某个主题的最后一条消息,新订阅者订阅后立即收到该消息 ### 4.3.4 遗嘱消息(Will Message)   客户端异常断开时,broker 自动发送预设消息 ### 4.3.5 清洁会话(Clean Session)   客户端断开连接后,broker 是否保留订阅信息和未送达消息 ## 4.4 组成   MQTT 消息由 固定头部(Fixed Header)、可变头部(Variable Header)、负载(Payload) 三部分组成: | 部分 | 作用 | |--------------|----------------------------------------------------------------------| | 固定头部 | 必选,2 字节起,包含消息类型(如发布、订阅)、QoS 等级、是否保留消息等标识。 | | 可变头部 | 可选,仅部分消息类型(如发布、订阅)需要,包含主题名、消息 ID 等信息。 | | 负载 | 可选,消息的实际内容(如 JSON 字符串、二进制数据),例如 {"temperature": 25}。 | # 五、.Net Core实现 ## 5.1 SuperSocket   [SuperSocket](https://www.supersocket.net/)是一个轻量级, 跨平台而且可扩展的 .Net/Mono Socket 服务器程序框架。可以轻松构建TCP、UDP、WebSocket服务器。 ## 5.2 安装依赖   NuGut安装SuperSocket、SuperSocket.WebSocket和SuperSocket.WebSocket.Server 2.0及以上版本。 ## 5.3 配置文件   appsettings.json ```json { "serverOptions": { "TcpServer": { "name": "TcpServer", "listeners": [ { "ip": "Any", "port": 4040 } ] }, "WebSocketServer": { "name": "WebSocket", "listeners": [ { "ip": "Any", "port": 5050 } ] } }, "Mqtt": { "Host": "127.0.0.1", "Port": 1883, "ClientId": "MqttClient", "Topics": [ "test/topic1", "test/topic2" ] } } ``` ## 5.3 主程序   program.cs ```cs var host = Host.CreateDefaultBuilder(args) .ConfigureServices((context, services) => { services.AddSingleton, TcpPackageHandler>(); services.AddSingleton(); // 1. 注入 MQTT 配置 services.Configure(context.Configuration.GetSection("Mqtt")); // 2. 注册 MQTT 客户端 services.AddSingleton(serviceProvider => new MqttClientFactory().CreateMqttClient()); // 3. 注册 MQTT 后台服务 services.AddHostedService(); }) .AsMultipleServerHostBuilder() .AddServer(builder => { builder.ConfigureServerOptions((ctx, config) => config.GetSection("TcpServer")); }) .AddWebSocketServer(builder => { builder .UseWebSocketMessageHandler(async (session, package) => { using var scope = session.Server.ServiceProvider.CreateScope(); var handler = scope.ServiceProvider.GetRequiredService(); await handler.HandleAsync(session, package); }) .ConfigureServerOptions((ctx, config) => config.GetSection("WebSocketServer")); }) .ConfigureLogging(logging => logging.ClearProviders()) .UseNLog() .Build(); await host.RunAsync(); ```   通过.AsMultipleServerHostBuilder()可以构造多服务器实例。 ## 5.4 TCP服务器 ### 5.4.1 协议数据 ```cs /// /// 协议数据结构 示例:65 6D 00 05 00 01 68 65 6C 6C 6F /// /// 帧头标识 /// 数据长度 /// 数据类型 /// 数据内容 /// record ProtocolFrame( ushort Magic, ushort Length, ushort Type, byte[] Payload ); ```   可以根据实际情况自定义消息格式。 ### 5.4.2 协议解析 ```cs class BinaryPipelineFilter : FixedHeaderPipelineFilter { /// /// 固定头长度 /// public BinaryPipelineFilter() : base(6) { } /// /// 从包头中解析出 Body 长度 /// /// 字节流 /// Body 长度 protected override int GetBodyLengthFromHeader(ref ReadOnlySequence buffer) { var reader = new SequenceReader(buffer); // 读取前2字节 → magic reader.TryReadBigEndian(out ushort magic); // 再读2字节 → Length reader.TryReadBigEndian(out ushort length); // 再读2字节 → Type reader.TryReadBigEndian(out ushort type); // 校验帧头 if (magic != 0x656D) { throw new Exception("非法帧头"); } // 限制长度 if (length == 0 || length > 8192) { throw new Exception("非法长度"); } // 返回 Payload 长度 return length; } /// /// 把完整字节包 → 转换成 ProtocolFrame /// /// 字节流 /// ProtocolFrame protected override ProtocolFrame DecodePackage(ref ReadOnlySequence buffer) { var reader = new SequenceReader(buffer); // 读取前2字节 → magic reader.TryReadBigEndian(out ushort magic); // 再读2字节 → Length reader.TryReadBigEndian(out ushort length); // 再读2字节 → Type reader.TryReadBigEndian(out ushort type); var payload = buffer.Slice(6, length).ToArray(); // 构造ProtocolFrame return new ProtocolFrame(magic, length, type, payload); } } ```   根据协议数据ProtocolFrame来解包。 ### 5.4.3 消息处理 ```cs class TcpPackageHandler(ILogger logger) : IPackageHandler { private readonly ILogger _logger = logger; public async ValueTask Handle(IAppSession session, ProtocolFrame package, CancellationToken cancellationToken) { _logger.LogInformation($"Magic={package.Magic:X4}, Type={package.Type}, Len={package.Length}"); await session.SendAsync(package.Payload, cancellationToken); } } ``` ### 5.4.4 TCPService ```cs class TcpService(IServiceProvider serviceProvider, IOptions serverOptions) : SuperSocketService(serviceProvider, serverOptions) { } ``` ## 5.5 Websocket服务器 ### 5.5.1 消息处理 ```cs class WebSocketMessageHandler(ILogger logger) { private readonly ILogger _logger = logger; public async ValueTask HandleAsync(WebSocketSession session, WebSocketPackage package) { _logger.LogInformation($"[WebSocket] {package.Message}"); await session.SendAsync("ok"); } } ``` ## 5.6 MQTT服务器 ### 5.6 MqttService ```cs class MqttService(IMqttClient mqttClient, IOptions settings, ILogger logger) : BackgroundService { private readonly IMqttClient _mqttClient = mqttClient; private readonly MqttSettings _settings = settings.Value; private readonly ILogger _logger = logger; private MqttClientOptions? _mqttOptions; protected override async Task ExecuteAsync(CancellationToken stoppingToken) { // 构建连接配置 _mqttOptions = new MqttClientOptionsBuilder() .WithClientId($"{_settings.ClientId}_{Guid.NewGuid():N}") .WithTcpServer(_settings.Host, _settings.Port) .WithCleanSession() .Build(); _mqttClient.ConnectedAsync += OnConnectedAsync; _mqttClient.DisconnectedAsync += OnDisconnectedAsync; _mqttClient.ApplicationMessageReceivedAsync += HandleMessage; _logger.LogInformation("[MQTT] 正在连接到服务器 {Host}:{Port}...", _settings.Host, _settings.Port); await _mqttClient.ConnectAsync(_mqttOptions, stoppingToken); // 等待程序停止 await Task.Delay(Timeout.Infinite, stoppingToken); } public async Task OnConnectedAsync(MqttClientConnectedEventArgs arg) { _logger.LogInformation("[MQTT] 已成功连接到服务器"); // 连接成功后订阅主题 foreach (var topic in _settings.Topics) { await _mqttClient.SubscribeAsync(topic, MqttQualityOfServiceLevel.AtLeastOnce); _logger.LogInformation($"[MQTT] 已订阅主题:{topic}"); } await Task.CompletedTask; } public async Task OnDisconnectedAsync(MqttClientDisconnectedEventArgs arg) { _logger.LogError("[MQTT] 连接失败,原因:{Reason}", arg.Reason); await Task.CompletedTask; } public async Task HandleMessage(MqttApplicationMessageReceivedEventArgs arg) { var topic = arg.ApplicationMessage.Topic; var payload = Encoding.UTF8.GetString(arg.ApplicationMessage.Payload); _logger.LogInformation($"\n[MQTT] 收到消息\n主题:{topic}\n内容:{payload}\n"); await Task.CompletedTask; } public async Task PublishAsync(string topic, string payload, CancellationToken cancellationToken = default) { try { var message = new MqttApplicationMessageBuilder() .WithTopic(topic) .WithPayload(payload) .WithQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce) .Build(); await _mqttClient.PublishAsync(message, cancellationToken); return true; } catch (Exception ex) { _logger.LogError(ex, "[MQTT] 发布消息到主题 {Topic} 失败", topic); return false; } } public override async Task StopAsync(CancellationToken stoppingToken) { await _mqttClient.DisconnectAsync(); await base.StopAsync(stoppingToken); } } ```   需要实现后台服务接口一直运行。