整体技术架构:
完整的数据处理流程可以概括为:
实时业务服务器
↓
MQTT Broker
↓
MQTT PUBLISH 消息
↓
WebSocket 二进制帧
↓
TLS 加密连接
↓
Python asyncio 客户端
↓
提取消息 Payload
↓
Base64 解码
↓
AES-CTR 解密
↓
自定义 TLV 协议解析
↓
转换为 Python 字典
↓
业务处理、存储或展示
程序首先使用 Python 的 asyncio 异步模型连接 MQTT 服务器。与普通 HTTP 轮询不同,MQTT 采用发布与订阅机制,客户端只需订阅指定 Topic,服务器有新数据时便会主动推送。这种方式延迟低、资源占用较少,适合实时比分、金融行情、设备状态、消息通知和监控系统等场景。
MQTT 数据通过 WebSocket 传输,连接地址使用 wss://,表示 WebSocket 外层启用了 TLS 加密。其完整协议链路可以理解为:MQTT 运行在 WebSocket 之上,WebSocket 运行在 TLS 和 TCP 之上。连接建立后,程序通过 Keep Alive 心跳维持在线状态,并可在网络异常时自动重连。
代码中的订阅等级为 QoS 0,特点是传输速度快、延迟较低,但不保证消息一定重传,因此适合高频且允许被后续数据覆盖的实时信息。
客户端接收到的 Payload 不一定是普通文本,也可能是 Base64 封装后的加密二进制数据。程序先进行 Base64 解码,再使用 AES-CTR 模式解密,随后按照自定义 TLV 格式解析字段。TLV 通常由字段名长度、字段名、字段值长度和字段值组成,比 JSON 更紧凑,适合网络实时传输。
代码还包含 FNV-1a 哈希、32 位整数运算和 JavaScript 无符号右移模拟,这类算法通常用于请求签名、字段校验、数据混淆或浏览器算法移植。
另外,WebRTC、ICE 和 STUN 部分主要用于收集本地地址、公网映射地址和网络候选信息。STUN 只能帮助发现 NAT 后的公网地址,并不等于 SOCKS5 代理。若需要代理,必须在底层网络连接中真正配置并使用代理服务器。
RESOURCES
文章资源
mqtt_websocket.py
import asyncio
from hbmqtt.client import MQTTClient
websocket_headers = {
'Accept-Language': 'zh-CN,zh;q=0.9,en;q=0.8,en-GB;q=0.7,en-US;q=0.6', 'Cache-Control': 'no-cache', 'Connection': 'Upgrade', 'Origin': 'https://widgets-livetracker.nami.com', 'Host': 'trackermq.namitiyu.com', 'Pragma': 'no-cache',
'Sec-WebSocket-Extensions': 'permessage-deflate; client_max_window_bits', 'Sec-WebSocket-Key':'8i6Yh1q7CbEAf13EKt5sBg==', 'Sec-WebSocket-Protocol': 'mqtt', 'Sec-WebSocket-Version': '13', 'Upgrade': 'websocket',
'User-Agent': 'Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/125.0.0.0 Safari/537.36'
}
websocket_url = 'wss://trackermq.namitiyu.com/mqtt'
proxy_uri = 'http://CC7C1AF8:75247E958880@tun-uzqqwl.qg.net:18031'
async def main(ID_):
client = MQTTClient()
await client.connect('wss://push.namitiyu.com/ws')
await client.subscribe([(f'live/m2/{ID_}', 0)])
while True:
message = await client.deliver_message()
print(message)
# asyncio.get_event_loop().run_until_complete(main(3800505))
import hbmqtt
import asyncio
from hbmqtt.client import MQTTClient, ClientException
from hbmqtt.mqtt.constants import QOS_1
print(hbmqtt)
# 配置 WebSocket 连接的参数
config = {
'keep_alive': 60, # 控制发送 ping 的间隔时间,以秒为单位
'ping_delay': 5, # 设置 ping 请求的延迟时间,以秒为单位
'auto_reconnect': True,
'reconnect_max_interval': 60,
'websockets': True, # 启用 WebSocket
'uri': websocket_url, # WebSocket 连接的 URI
'extra_headers': websocket_headers
}
async def connect_and_subscribe(ID_):
client = MQTTClient(config=config)
try:
await client.connect(config['uri'])
print("Connected to MQTT broker over WebSocket")
# 订阅一个主题
await client.subscribe([(f'live/m2/{ID_}', 0)])
print("Subscribed to 'test/topic'")
while True:
# 等待消息
message = await client.deliver_message()
packet = message.publish_packet
print(f"Received message: {packet.payload.data.decode()}")
except ClientException as ce:
print(f"Client exception: {ce}")
finally:
await client.disconnect()
loop = asyncio.get_event_loop()
loop.run_until_complete(connect_and_subscribe(4400588))
评论交流
还没有公开评论。