feat:更新文档架构

This commit is contained in:
2026-05-20 11:26:38 +08:00
parent ad9a4d6806
commit 8c9b3adc57
46 changed files with 218 additions and 75 deletions

View File

@@ -0,0 +1,943 @@
---
title: 网络编程简介
date: 2025-12-15
---
# 一、TCP
## 1.1 定义
  TCPTransmission 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 定义
  HTTPHyperText 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.01996 年)
- 基础版本,支持 GET、POST、HEAD 三种请求方法。
- 每次请求都需要建立新的 TCP 连接(短连接),连接建立和关闭的开销大,性能较低。
- 不支持长连接、管线化请求,仅支持简单的文本传输。
### 2.3.2 HTTP/1.11999 年)
- 默认开启 长连接Keep-Alive同一 TCP 连接可处理多个请求,减少连接开销。
- 支持 管线化请求:客户端可连续发送多个请求,无需等待前一个响应返回(部分浏览器未完全支持)。
- 新增请求方法PUT、DELETE、OPTIONS、TRACE、CONNECT。
- 支持 chunked 编码分块传输、缓存控制Cache-Control、内容协商等核心功能。
### 2.3.3 HTTP/22015 年)
  基于 SPDY 协议优化,核心目标是提升性能:
- 二进制帧传输:将请求 / 响应数据拆分为二进制帧,而非 HTTP/1.x 的文本格式,解析效率更高。
- 多路复用:同一 TCP 连接中可并发处理多个请求(通过帧的 Stream ID 区分),解决 HTTP/1.1 的 “队头阻塞” 问题。
- 服务器推送Server Push服务器可主动向客户端推送关联资源如 HTML 引用的 CSS/JS减少客户端请求次数。
- 头部压缩HPACK对请求头和响应头进行压缩减少传输体积HTTP/1.x 头部重复传输开销大)。
### 2.3.4 HTTP/32022 年)
  基于 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 <token>
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。
#### 请求头
&emsp;&emsp;描述请求的附加信息,常用字段:
- Host目标服务器域名如 lms.example.comHTTP/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 表示不使用缓存)。
#### 请求体
&emsp;&emsp;可选部分,仅在需要向服务器提交数据时使用(如 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
&emsp;&emsp;服务器向客户端返回的响应格式,由 状态行、响应头、空行、响应体 四部分组成:
```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网关超时|
#### 响应头
&emsp;&emsp;描述响应的附加信息,常用字段:
- 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 状态码必选)。
#### 响应体
&emsp;&emsp;服务器返回的核心数据,数据格式由 Content-Type 指定,常见格式:
- JSON前后端分离项目首选如 {"code": 200, "data": [...]})。
- HTML传统 Web 页面如静态网页、JSP 页面)。
- 图片 / 视频:二进制流(如 image/jpeg、video/mp4
- 纯文本text/plain。
# 三、WebSocket
## 3.1 定义
&emsp;&emsp;WebSocket 是一种**全双工、双向、持久化的网络通信协议**(属于应用层协议),由 HTML5 规范定义,专门解决 HTTP 协议无法实现服务器主动向客户端推送数据的问题。
&emsp;&emsp;它的核心特点是:**一旦客户端与服务器建立连接,双方就可以在这个连接上实时、双向地发送数据**,无需像 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 帧格式(二进制 / 文本)。
&emsp;&emsp;请求头:
```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 # 跨域验证
```
&emsp;&emsp;响应头:
```http
HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo= # 服务器加密后的密钥,客户端验证
```
### 3.2.2 全双工通信
&emsp;&emsp;此时 HTTP 连接已升级为 WebSocket 连接双方可以通过这个连接实时、双向地发送数据。数据传输采用帧Frame 格式,支持文本数据和二进制数据(如图片、视频流)。
### 3.2.3 帧格式
&emsp;&emsp;WebSocket 数据以帧为单位传输,帧头包含操作码(文本帧 0x01、二进制帧 0x02、关闭帧 0x08 等)、掩码(客户端发送数据需掩码,服务器无需);
### 3.2.4 无同源限制
&emsp;&emsp;WebSocket 不遵循同源策略(但服务器可通过 Origin 头限制跨域);
### 3.2.5 心跳机制
&emsp;&emsp;通过 Ping/Pong 帧维持连接(避免网络设备断开空闲连接,如 LMS 系统需定期发送 Ping 帧,服务器回复 Pong 帧)。
# 四、MQTT
## 4.1 定义
&emsp;&emsp;MQTTMessage 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 轻量级
&emsp;&emsp;头部开销极小(固定头部仅 2 字节),消息体支持二进制 / 文本,适合低带宽场景
### 4.3.3 保留消息Retained Message
&emsp;&emsp;broker 保存某个主题的最后一条消息,新订阅者订阅后立即收到该消息
### 4.3.4 遗嘱消息Will Message
&emsp;&emsp;客户端异常断开时broker 自动发送预设消息
### 4.3.5 清洁会话Clean Session
&emsp;&emsp;客户端断开连接后broker 是否保留订阅信息和未送达消息
## 4.4 组成
&emsp;&emsp;MQTT 消息由 固定头部Fixed Header、可变头部Variable Header、负载Payload 三部分组成:
| 部分 | 作用 |
|--------------|----------------------------------------------------------------------|
| 固定头部 | 必选2 字节起包含消息类型如发布、订阅、QoS 等级、是否保留消息等标识。 |
| 可变头部 | 可选,仅部分消息类型(如发布、订阅)需要,包含主题名、消息 ID 等信息。 |
| 负载 | 可选,消息的实际内容(如 JSON 字符串、二进制数据),例如 {"temperature": 25}。 |
# 五、.Net Core实现
## 5.1 SuperSocket
&emsp;&emsp;[SuperSocket](https://www.supersocket.net/)是一个轻量级, 跨平台而且可扩展的 .Net/Mono Socket 服务器程序框架。可以轻松构建TCP、UDP、WebSocket服务器。
## 5.2 安装依赖
&emsp;&emsp;NuGut安装SuperSocket、SuperSocket.WebSocket和SuperSocket.WebSocket.Server 2.0及以上版本。
## 5.3 配置文件
&emsp;&emsp;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 主程序
&emsp;&emsp;program.cs
```cs
var host = Host.CreateDefaultBuilder(args)
.ConfigureServices((context, services) =>
{
services.AddSingleton<IPackageHandler<ProtocolFrame>, TcpPackageHandler>();
services.AddSingleton<WebSocketMessageHandler>();
// 1. 注入 MQTT 配置
services.Configure<MqttSettings>(context.Configuration.GetSection("Mqtt"));
// 2. 注册 MQTT 客户端
services.AddSingleton<IMqttClient>(serviceProvider => new MqttClientFactory().CreateMqttClient());
// 3. 注册 MQTT 后台服务
services.AddHostedService<MqttService>();
})
.AsMultipleServerHostBuilder()
.AddServer<TcpService, ProtocolFrame, BinaryPipelineFilter>(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<WebSocketMessageHandler>();
await handler.HandleAsync(session, package);
})
.ConfigureServerOptions((ctx, config) => config.GetSection("WebSocketServer"));
})
.ConfigureLogging(logging => logging.ClearProviders())
.UseNLog()
.Build();
await host.RunAsync();
```
&emsp;&emsp;通过.AsMultipleServerHostBuilder()可以构造多服务器实例。
## 5.4 TCP服务器
### 5.4.1 协议数据
```cs
/// <summary>
/// 协议数据结构 示例65 6D 00 05 00 01 68 65 6C 6C 6F
/// </summary>
/// <param name="Magic">帧头标识</param>
/// <param name="Length">数据长度</param>
/// <param name="Type">数据类型</param>
/// <param name="Payload">数据内容</param>
///
record ProtocolFrame(
ushort Magic,
ushort Length,
ushort Type,
byte[] Payload
);
```
&emsp;&emsp;可以根据实际情况自定义消息格式。
### 5.4.2 协议解析
```cs
class BinaryPipelineFilter : FixedHeaderPipelineFilter<ProtocolFrame>
{
/// <summary>
/// 固定头长度
/// </summary>
public BinaryPipelineFilter() : base(6)
{
}
/// <summary>
/// 从包头中解析出 Body 长度
/// </summary>
/// <param name="buffer">字节流</param>
/// <returns>Body 长度</returns>
protected override int GetBodyLengthFromHeader(ref ReadOnlySequence<byte> buffer)
{
var reader = new SequenceReader<byte>(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;
}
/// <summary>
/// 把完整字节包 → 转换成 ProtocolFrame
/// </summary>
/// <param name="buffer">字节流</param>
/// <returns>ProtocolFrame</returns>
protected override ProtocolFrame DecodePackage(ref ReadOnlySequence<byte> buffer)
{
var reader = new SequenceReader<byte>(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);
}
}
```
&emsp;&emsp;根据协议数据ProtocolFrame来解包。
### 5.4.3 消息处理
```cs
class TcpPackageHandler(ILogger<TcpPackageHandler> logger) : IPackageHandler<ProtocolFrame>
{
private readonly ILogger<TcpPackageHandler> _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> serverOptions) : SuperSocketService<ProtocolFrame>(serviceProvider, serverOptions)
{
}
```
## 5.5 Websocket服务器
### 5.5.1 消息处理
```cs
class WebSocketMessageHandler(ILogger<WebSocketMessageHandler> logger)
{
private readonly ILogger<WebSocketMessageHandler> _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<MqttSettings> settings, ILogger<TcpPackageHandler> logger) : BackgroundService
{
private readonly IMqttClient _mqttClient = mqttClient;
private readonly MqttSettings _settings = settings.Value;
private readonly ILogger<TcpPackageHandler> _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<bool> 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);
}
}
```
&emsp;&emsp;需要实现后台服务接口一直运行。