
标签:#物联网 #Modbus #TCP/IP #UDP #POE供电 #腾讯云 #Wireshark #Python #InfluxDB #以太网温湿度传感器 #网口温湿度记录仪 #机房监控
上一篇我们做了 TCP vs UDP 的基准测试,结论很明确:在跨城多机房、高并发轮询场景下,中心平台直接拉取数百个 Modbus UDP 节点既不经济也不可靠。更合理的架构是让传感器主动上报,平台侧只负责监听。
但问题来了:当传感器以 UDP 广播/单播方式主动推送数据时,标准 socket 监听会遇到几个现实障碍:
这时候,Raw Socket 就是唯一的选择——直接让内核把 IP 层(甚至以太网层)的原始报文递交给用户态,绕过 TCP/UDP 协议栈的应用层绑定逻辑。
模式 | Socket 类型 | 能拿到什么 | 适用场景 |
|---|---|---|---|
IP Raw Socket | SOCK_RAW + IPPROTO_UDP | IP 头 + UDP 头 + 载荷 | 监听特定协议的所有报文,内核帮你剥掉链路层 |
Packet Socket | SOCK_RAW + ETH_P_ALL | 以太网头 + IP 头 + UDP 头 + 载荷 | 需要 MAC 地址、VLAN tag 等链路层信息 |
AF_PACKET | SOCK_RAW + htons(ETH_P_ALL) | 同 Packet Socket,Linux 特有高性能接口 | 高吞吐抓包、自定义过滤 |
本文重点使用 IP Raw Socket(模式一),兼顾可用性和解析深度。需要链路层信息时,再切到 Packet Socket。
以常见的工业以太网温湿度记录仪为例,其主动上报的 UDP 报文结构如下:
┌─────────────────────────────────────────────────────────────┐
│ 以太网帧头 (14 bytes) │
│ DST MAC (6) | SRC MAC (6) | EtherType (2) = 0x0800 │
├─────────────────────────────────────────────────────────────┤
│ IP 头 (20 bytes, 不含选项) │
│ Ver/IHL | TOS | Total Len | ID | Flags/Offset │
│ TTL | Protocol (17=UDP) | Header Checksum │
│ Src IP (4) | Dst IP (4) │
├─────────────────────────────────────────────────────────────┤
│ UDP 头 (8 bytes) │
│ Src Port (2) | Dst Port (2) = 502 │
│ Length (2) | Checksum (2) │
├─────────────────────────────────────────────────────────────┤
│ Modbus TCP ADU (变长) │
│ Transaction ID (2) | Protocol ID (2) = 0x0000 │
│ Length (2) | Unit ID (1) │
├─────────────────────────────────────────────────────────────┤
│ Modbus PDU (变长) │
│ Function Code (1) = 0x04 (读输入寄存器) │
│ Byte Count (1) │
│ Register Values (N × 2 bytes) │
│ ├── 寄存器 0: 温度 ×10, INT16 │
│ ├── 寄存器 1: 湿度 ×10, INT16 │
│ └── 寄存器 2: 露点 ×10, INT16 │
└─────────────────────────────────────────────────────────────┘#!/usr/bin/env python3
"""
raw_udp_listener.py - 基于 Raw Socket 的 UDP 报文监听与 Modbus 负载解析
"""
import socket
import struct
import binascii
import logging
from datetime import datetime
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s [%(levelname)s] %(message)s'
)
# ---------- 全局配置 ----------
LISTEN_IP = "0.0.0.0" # 监听所有接口
TARGET_PORT = 502 # Modbus 默认端口
BUFFER_SIZE = 65535 # 最大 IP 包大小
def create_raw_socket():
"""
创建 IP Raw Socket,监听 UDP 协议。
注意:
- Linux 需要 CAP_NET_RAW 权限(root 或 setcap)
- IP_HDRINCL 告诉内核我们收到的报文包含 IP 头
"""
try:
sock = socket.socket(
socket.AF_INET, # IPv4
socket.SOCK_RAW, # Raw Socket
socket.IPPROTO_UDP # 只接收 UDP 报文
)
except PermissionError:
logging.error("需要 root 权限或 CAP_NET_RAW 能力")
raise
except socket.error as e:
logging.error(f"创建 Raw Socket 失败: {e}")
raise
# 设置 IP_HDRINCL,让 recvfrom 返回包含 IP 头的完整报文
sock.setsockopt(socket.IPPROTO_IP, socket.IP_HDRINCL, 1)
# 设置接收缓冲区(高吞吐场景)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 8 * 1024 * 1024)
# 绑定(Raw Socket 绑定 IP 而非端口)
sock.bind((LISTEN_IP, 0))
logging.info(f"Raw Socket 已创建,监听 UDP 流量 (端口过滤在应用层完成)")
return sockdef parse_ip_header(data: bytes) -> tuple:
"""
解析 IPv4 头部,返回 (src_ip, dst_ip, protocol, header_len, total_len)
IP 头结构(前 20 字节):
Byte 0: Ver(4) + IHL(4)
Byte 1: TOS
Byte 2-3: Total Length (16 bits, big-endian)
Byte 4-5: Identification
Byte 6-7: Flags(3) + Fragment Offset(13)
Byte 8: TTL
Byte 9: Protocol (17=UDP)
Byte 10-11: Header Checksum
Byte 12-15: Source IP
Byte 16-19: Destination IP
"""
if len(data) < 20:
return None
ver_ihl = data[0]
version = ver_ihl >> 4
ihl = ver_ihl & 0x0F
header_len = ihl * 4 # IHL 单位是 32-bit word
total_len = struct.unpack(">H", data[2:4])[0]
protocol = data[9]
src_ip = socket.inet_ntoa(data[12:16])
dst_ip = socket.inet_ntoa(data[16:20])
return src_ip, dst_ip, protocol, header_len, total_lendef parse_udp_header(data: bytes, ip_header_len: int) -> tuple:
"""
解析 UDP 头部,返回 (src_port, dst_port, udp_len, udp_checksum, payload)
UDP 头结构(8 字节):
Byte 0-1: Source Port
Byte 2-3: Destination Port
Byte 4-5: Length (UDP 头 + 载荷)
Byte 6-7: Checksum
"""
udp_offset = ip_header_len
if len(data) < udp_offset + 8:
return None
src_port, dst_port, udp_len, udp_checksum = struct.unpack(
">HHHH", data[udp_offset:udp_offset + 8]
)
payload_offset = udp_offset + 8
payload = data[payload_offset:]
return src_port, dst_port, udp_len, udp_checksum, payloaddef parse_modbus_adu(payload: bytes) -> dict | None:
"""
解析 Modbus TCP ADU (7 字节 MBAP + PDU)
MBAP:
- Transaction ID (2 bytes)
- Protocol ID (2 bytes, 0x0000)
- Length (2 bytes, 后续字节数)
- Unit ID (1 byte)
PDU (Function Code + Data):
- 0x03: 读保持寄存器
- 0x04: 读输入寄存器
"""
if len(payload) < 8:
return None
tid, proto_id, length, unit_id = struct.unpack(">HHHB", payload[:7])
pdu = payload[7:]
if len(pdu) < 2:
return None
func_code = pdu[0]
# 异常响应检测
if func_code & 0x80:
exception_code = pdu[1] if len(pdu) > 1 else 0
logging.warning(f"Modbus 异常: FC=0x{func_code:02X}, Exception=0x{exception_code:02X}")
return None
result = {
"transaction_id": tid,
"protocol_id": proto_id,
"length": length,
"unit_id": unit_id,
"function_code": func_code,
"registers": [],
}
# 解析寄存器数据
if func_code in (0x03, 0x04): # 读保持/输入寄存器
if len(pdu) < 2:
return result
byte_count = pdu[1]
reg_data = pdu[2:2 + byte_count]
regs = []
for i in range(0, len(reg_data), 2):
val = struct.unpack(">H", reg_data[i:i+2])[0]
# 处理有符号(温度可能为负)
if val > 32767:
val -= 65536
regs.append(val)
result["registers"] = regs
# 转换为工程值
if len(regs) >= 2:
result["temperature"] = round(regs[0] / 10.0, 1)
result["humidity"] = round(regs[1] / 10.0, 1)
if len(regs) >= 3:
dew = regs[2]
result["dew_point"] = round(dew / 10.0, 1)
return resultdef main():
sock = create_raw_socket()
packet_count = 0
modbus_count = 0
logging.info("开始监听... (Ctrl+C 停止)")
try:
while True:
raw_data, addr = sock.recvfrom(BUFFER_SIZE)
packet_count += 1
# 1. 解析 IP 头
ip_info = parse_ip_header(raw_data)
if not ip_info:
continue
src_ip, dst_ip, protocol, ip_hdr_len, total_len = ip_info
# 只处理 UDP
if protocol != 17:
continue
# 2. 解析 UDP 头
udp_info = parse_udp_header(raw_data, ip_hdr_len)
if not udp_info:
continue
src_port, dst_port, udp_len, udp_csum, payload = udp_info
# 3. 端口过滤(应用层过滤,Raw Socket 不帮你做)
if dst_port != TARGET_PORT and src_port != TARGET_PORT:
continue
# 4. 解析 Modbus 负载
modbus_data = parse_modbus_adu(payload)
if not modbus_data:
continue
modbus_count += 1
# 5. 输出结果
ts = datetime.now().strftime("%H:%M:%S.%f")[:-3]
print(f"[{ts}] {src_ip}:{src_port} → {dst_ip}:{dst_port}")
print(f" Transaction ID: {modbus_data['transaction_id']}")
print(f" Unit ID: {modbus_data['unit_id']}")
print(f" Function: 0x{modbus_data['function_code']:02X}")
if "temperature" in modbus_data:
print(f" 🌡️ 温度: {modbus_data['temperature']}℃")
print(f" 💧 湿度: {modbus_data['humidity']}%RH")
if "dew_point" in modbus_data:
print(f" 💧 露点: {modbus_data['dew_point']}℃")
else:
print(f" 📊 原始寄存器: {modbus_data['registers']}")
print()
# 6. 写入 InfluxDB / MQTT(见前篇)
# publish_to_influx(modbus_data, src_ip)
except KeyboardInterrupt:
logging.info(f"停止监听. 总报文: {packet_count}, Modbus 报文: {modbus_count}")
finally:
sock.close()
if __name__ == "__main__":
main()# 方法一:直接用 root 运行(简单但不推荐生产)
sudo python3 raw_udp_listener.py
# 方法二:给 Python 解释器赋予 CAP_NET_RAW 能力(推荐)
sudo setcap cap_net_raw+eip /usr/bin/python3
# 验证
getcap /usr/bin/python3
# 输出: /usr/bin/python3 = cap_net_raw+eip
# 方法三:用 systemd 服务,配置 CapabilityBoundingSet# /etc/systemd/system/raw-sensor-listener.service
[Unit]
Description=Raw Socket UDP Listener for Ethernet Temperature Sensors
After=network.target
[Service]
Type=simple
User=sensor
Group=sensor
ExecStart=/usr/bin/python3 /opt/sensor-listener/raw_udp_listener.py
Restart=always
RestartSec=5
# 能力配置
CapabilityBoundingSet=CAP_NET_RAW
AmbientCapabilities=CAP_NET_RAW
# 安全加固
NoNewPrivileges=true
PrivateTmp=true
ProtectSystem=strict
ReadWritePaths=/var/log/sensor-listener
StandardOutput=journal
StandardError=journal
[Install]
WantedBy=multi-user.targetRaw Socket 默认把所有 UDP 报文都递交给用户态,开销巨大。可以用 SO_ATTACH_FILTER 挂载 BPF 过滤器,让内核只把目标端口为 502 的报文递上来:
# BPF 过滤: UDP 目标端口 == 502
# 以太网头 14 字节, IP 头 20 字节, UDP 头偏移 2 字节是 dst_port
# 指令序列:
# 加载 UDP 目标端口 (以太网14 + IP20 + UDP2 = 偏移36)
# 比较是否等于 502 (0x01F6)
BPF_FILTER = [
0x28, 0x00, 0x00, 0x0000000C, # ldh [12] - 加载 EtherType
0x15, 0x00, 0x05, 0x00000800, # jeq #0x800 (IP) goto next, else jump 5
0x30, 0x00, 0x00, 0x00000017, # ldb [23] - 加载 IP protocol
0x15, 0x00, 0x03, 0x00000011, # jeq #17 (UDP) goto next, else jump 3
0x28, 0x00, 0x00, 0x00000024, # ldh [36] - 加载 UDP dst_port
0x15, 0x00, 0x01, 0x000001F6, # jeq #502 goto next, else jump 1
0x06, 0x00, 0x00, 0x0000FFFF, # ret #0xFFFF (接受)
0x06, 0x00, 0x00, 0x00000000, # ret #0x0 (拒绝)
]
import ctypes
BpfProg = struct.pack("HL", len(BPF_FILTER), ctypes.addressof(
(ctypes.c_uint32 * len(BPF_FILTER))(*BPF_FILTER)
))
# 实际使用中建议用 python-bcc 或 pyroute2 库来设置 BPF更简单的做法:用 libpcap/scapy 的 BPF 过滤字符串 "udp dst port 502",内核自动编译为 BPF 字节码。
import multiprocessing as mp
def worker(queue, worker_id):
"""每个 worker 进程一个 Raw Socket"""
sock = create_raw_socket()
# 设置 BPF 过滤
# ...
while True:
data, addr = sock.recvfrom(BUFFER_SIZE)
queue.put((data, addr))
def collector(queue):
"""汇总进程:解析 + 入库"""
while True:
data, addr = queue.get()
result = parse_pipeline(data)
if result:
write_to_influx(result)
def main_mp():
queue = mp.Queue(maxsize=10000)
workers = []
for i in range(4): # 4 个 CPU 核心
p = mp.Process(target=worker, args=(queue, i))
p.start()
workers.append(p)
c = mp.Process(target=collector, args=(queue,))
c.start()
c.join()方案 | 性能 | 复杂度 | 灵活性 | 适用场景 |
|---|---|---|---|---|
Raw Socket (SOCK_RAW) | 中 | 低 | 中 | 通用监听、端口过滤 |
AF_PACKET + TPACKETv3 | 高 | 中 | 高 | 高吞吐、零拷贝 |
eBPF/XDP | 极高 | 高 | 极高 | 10G+ 链路、内核态过滤 |
libpcap/scapy | 低~中 | 极低 | 低 | 开发调试、原型验证 |
Raw Socket 拿到的数据和 Wireshark 抓到的数据应该完全一致。交叉验证方法:
# 1. 用 tshark 同时抓包
tshark -i eth0 -f "udp port 502" -w /tmp/capture.pcap &
# 2. 运行你的 Raw Socket 程序
python3 raw_udp_listener.py &
# 3. 对比计数
# Wireshark 显示 10000 个 UDP/502 报文
# 你的程序也显示 10000 个 Modbus 报文 → 一致 ✓
# 4. 用 Wireshark 解码验证解析正确性
tshark -r /tmp/capture.pcap -Y "modbus" -T fields \
-e ip.src -e udp.srcport -e modbus.reg16Raw Socket 能接收网卡上的所有 UDP 报文,生产环境务必做好隔离:
# 只允许来自传感器子网的报文
iptables -A INPUT -p udp --dport 502 -s 192.168.10.0/24 -j ACCEPT
iptables -A INPUT -p udp --dport 502 -j DROPRaw Socket 绕过了内核协议栈的校验,应用层必须自己做:
import zlib
def verify_udp_checksum(ip_header, udp_header_and_payload):
"""验证 UDP 校验和(可选,多数局域网场景可跳过)"""
# 构造伪头部
src_ip = ip_header[12:16]
dst_ip = ip_header[16:20]
zero = b'\x00'
protocol = b'\x11' # UDP
udp_len = struct.pack('>H', len(udp_header_and_payload))
pseudo_header = src_ip + dst_ip + zero + protocol + udp_len
checksum_data = pseudo_header + udp_header_and_payload
# 补齐奇数长度
if len(checksum_data) % 2:
checksum_data += b'\x00'
# 计算 16-bit one's complement sum
total = 0
for i in range(0, len(checksum_data), 2):
total += (checksum_data[i] << 8) + checksum_data[i+1]
while total >> 16:
total = (total & 0xFFFF) + (total >> 16)
checksum = ~total & 0xFFFF
return checksum == 0struct.unpack(">H", ...) 时注意 > 前缀。 scapy 的 L2Socket + sniff 做快速原型 基于 Raw Socket 监听以太网温湿度传感器的 UDP 报文,核心价值在于:
但这套方案不适合作为生产采集的主路径。它的正确定位是:
调试工具 + 旁路监控 + 流量分析
生产环境的主路径仍然是:传感器 → 标准 UDP/TCP → 边缘网关(Modbus 协议栈)→ MQTT → 云端。Raw Socket 监听作为旁路,用于验证数据完整性、排查丢包、分析协议异常。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。