feat:增加博客统计组件
This commit is contained in:
@@ -1,36 +1,38 @@
|
||||
# TCP
|
||||
## 定义
|
||||
  TCP(Transmission Control Protocol,传输控制协议)是互联网核心的面向连接、可靠、字节流的传输层协议,工作在 OSI 模型的传输层(TCP/IP 模型的传输层),基于 IP 协议提供端到端的可靠数据传输服务。它是 HTTP、HTTPS、WebSocket、MQTT 等应用层协议的底层依赖,核心目标是解决 IP 协议 “无连接、不可靠、无顺序” 的缺陷,确保数据在不可靠的网络中准确、完整、有序地传输。
|
||||
<ArticleMetadata />
|
||||
|
||||
## 特性
|
||||
### 面向连接
|
||||
# 一、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. 支持流量 / 拥塞控制
|
||||
@@ -41,17 +43,404 @@
|
||||
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
|
||||
## 定义
|
||||
  HTTP(HyperText Transfer Protocol,超文本传输协议)是互联网的核心协议之一,用于客户端(如浏览器、App)与服务器之间的通信,是万维网(WWW)数据交换的基础。它定义了请求 / 响应的格式、传输规则和状态码等核心机制,支持从简单文本到复杂多媒体(图片、视频、文件)的传输,也是现代 Web 应用的底层通信标准。
|
||||
  HTTP(HyperText Transfer Protocol,超文本传输协议)是互联网的核心协议之一,用于**客户端(如浏览器、App)与服务器之间的通信**,是万维网(WWW)数据交换的基础。它定义了请求 / 响应的格式、传输规则和状态码等核心机制,支持从简单文本到复杂多媒体(图片、视频、文件)的传输,也是现代 Web 应用的底层通信标准。
|
||||
|
||||
## 特性
|
||||
### 请求 - 响应模式
|
||||
  通信由客户端主动发起请求,服务器接收后处理并返回响应,不存在服务器主动向客户端推送数据的情况(HTTP/2 引入 Server Push 扩展,可主动推送关联资源)。
|
||||
  一次完整通信流程:客户端建立连接 → 发送请求 → 服务器处理 → 返回响应 → 连接关闭(HTTP/1.1 默认开启长连接 Keep-Alive)。
|
||||
  一次完整通信流程:**客户端建立连接 → 发送请求 → 服务器处理 → 返回响应 → 连接关闭**(HTTP/1.1 默认开启长连接 Keep-Alive)。
|
||||
|
||||
### 无状态
|
||||
  服务器不会保存客户端的会话状态(如登录状态、浏览记录),每次请求都是独立的,服务器无法通过协议本身识别连续请求是否来自同一客户端。通过 Cookie、Session、Token(如 JWT)等机制补充状态管理。
|
||||
  **服务器不会保存客户端的会话状态**(如登录状态、浏览记录),每次请求都是独立的,服务器无法通过协议本身识别连续请求是否来自同一客户端。通过 Cookie、Session、Token(如 JWT)等机制补充状态管理。
|
||||
|
||||
## 版本
|
||||
### HTTP/1.0(1996 年)
|
||||
@@ -159,8 +548,8 @@ Set-Cookie: sessionId=abc123; Path=/
|
||||
|
||||
# WebSocket
|
||||
## 定义
|
||||
  它的核心特点是:一旦客户端与服务器建立连接,双方就可以在这个连接上实时、双向地发送数据,无需像 HTTP 那样每次通信都由客户端发起请求,非常适合实时通信场景(如聊天、直播弹幕、实时数据监控、在线协作等)。
|
||||
|
||||
  WebSocket 是一种**全双工、双向、持久化的网络通信协议**(属于应用层协议),由 HTML5 规范定义,专门解决 HTTP 协议无法实现服务器主动向客户端推送数据的问题。
|
||||
  它的核心特点是:**一旦客户端与服务器建立连接,双方就可以在这个连接上实时、双向地发送数据**,无需像 HTTP 那样每次通信都由客户端发起请求,非常适合实时通信场景(如聊天、直播弹幕、实时数据监控、在线协作等)。
|
||||
|
||||
| 特性 | HTTP | WebSocket |
|
||||
|---------------------|-------------------------------|-------------------------------|
|
||||
@@ -209,7 +598,7 @@ Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo= # 服务器加密后的密
|
||||
|
||||
# MQTT
|
||||
## 定义
|
||||
|
||||
  MQTT(Message Queuing Telemetry Transport,消息队列遥测传输)是一种**轻量级、低带宽、低功耗的发布 / 订阅(Publish/Subscribe)模式物联网(IoT)通信协议**,由 IBM 于 1999 年设计,核心目标是解决受限设备(如传感器、嵌入式设备)和低带宽、不稳定网络环境下的高效数据传输问题。
|
||||
|
||||
## 架构
|
||||
- 发布者(Publisher):发送消息的设备 / 服务;
|
||||
|
||||
Reference in New Issue
Block a user