ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Python asyncio 高并发 Modbus TCP 采集实战:200台设备秒级轮询

Python asyncio 高并发 Modbus TCP 采集实战:200台设备秒级轮询 1. 项目缘起与整体架构思路1.1 为什么会有这个采集需求做过机房动环监控或者仓储环境监测的朋友应该都有体会当温湿度探头数量从几台涨到几十台、上百台之后传统的轮询方式就开始力不从心了。我之前接手的一个项目现场有 200 多台支持 POE 供电的以太网温湿度变送器分布在三层楼、十几个弱电间里。这些设备全部走 Modbus TCP/IP 协议每台设备一个 IP端口统一是 502。最初的方案是用一个for循环挨个socket连接、发请求、读响应、断开一轮下来差不多要 40 多秒。问题是这些变送器的数据刷新周期本身就只有几秒等一轮轮询完最早读到的数据已经过期了。更麻烦的是只要有一台设备网络抖动或者掉线整个循环就会被socket.timeout卡住后面的设备全部排队等着实时性彻底崩掉。后来我把方案换成了asyncio协程并发同样的 200 多台设备一轮采集稳定在 1.5 秒以内单台设备超时也不会影响其他设备。这篇文章就把这套方案的完整思路、代码实现和踩过的坑整理出来适合有一定 Python 基础、正在做工业数据采集或者物联网网关开发的朋友参考。哪怕你之前没怎么用过 asyncio跟着思路走也能理解为什么这么设计。1.2 整体架构从串行到并发的核心转变先说清楚这个项目的核心矛盾Modbus TCP 是请求-响应式的同步协议而我们要同时和几百台设备打交道。这两件事天然是冲突的。Modbus 协议本身很简单一个请求报文发出去设备回一个响应报文一来一回。它没有推送机制你想拿数据就必须主动去问。所以同时轮询数百台设备本质上不是真的同时而是在极短的时间窗口内把几百个问-答过程并发地铺开让它们的时间互相重叠而不是排队。这里就要区分两个概念并发concurrency和并行parallelism。并行是真的一起跑需要多核 CPU并发是看起来一起跑靠的是在等待 IO 的时候切换去干别的事。Modbus TCP 采集是典型的IO 密集型任务——绝大部分时间都花在等网络响应上CPU 几乎不干活。这种场景下asyncio 的单线程事件循环就是最优解开多线程反而会因为 GIL 和线程切换带来额外开销。我的整体架构是这样的设备清单一个配置文件或者数据库表存所有设备的 IP、端口、从站地址、寄存器地址映射。协程池用asyncio.Semaphore控制并发上限避免一次性把几百个连接全砸出去把交换机打爆。单设备采集协程每个协程负责一台设备的完整连接-请求-解析-关闭流程。超时与重试每台设备独立超时失败重试绝不拖累别人。数据汇聚所有协程的结果通过asyncio.gather收集统一入库或上报。这个架构的关键在于隔离性每台设备的失败被限制在自己的协程里不会像串行方案那样一颗老鼠屎坏了一锅粥。1.3 为什么不用现成的 Modbus 库市面上有pymodbus这样的成熟库功能很全。但在高并发场景下我最终选择了自己用 asyncio 裸写 Modbus TCP 报文原因有几个第一pymodbus的同步客户端在并发场景下要么开线程池要么用它的异步客户端但异步客户端的连接管理在高频短连接场景下开销不小。第二我们的需求非常单一——只读保持寄存器功能码 0x03和输入寄存器功能码 0x04不需要写、不需要复杂的事务管理。第三自己控制报文意味着可以精确控制超时、可以复用连接、可以做批量读取优化。当然如果你只是偶尔读几台设备用pymodbus完全没问题没必要重复造轮子。但当你面对几百台设备、要求秒级刷新的时候理解底层报文结构、自己掌控 IO 就变得很有价值了。下面我会把 Modbus TCP 的报文结构讲清楚这样你自己写或者调试pymodbus都用得上。2. Modbus TCP 报文结构与采集核心细节2.1 Modbus TCP 报文到底长什么样很多人用pymodbus用久了反而说不清楚报文里每个字节是什么。我建议你一定要把这部分搞明白因为现场调试的时候抓包工具里看到的是一串十六进制看不懂报文就只能干瞪眼。Modbus TCP 的报文叫MBAP 头Modbus Application Protocol header PDU协议数据单元。MBAP 头固定 7 个字节PDU 长度可变。MBAP 头的结构如下字段长度说明事务标识符 Transaction ID2 字节请求和响应配对用自己生成一般递增协议标识符 Protocol ID2 字节Modbus TCP 固定为 0x0000长度 Length2 字节后面还有多少字节单元标识符 PDU单元标识符 Unit ID1 字节从站地址TCP 场景下常用来区分网关后的串口设备PDU 部分读保持寄存器的请求是功能码 0x03 起始地址2 字节 寄存器数量2 字节。响应是功能码 0x03 字节数1 字节 数据N 字节。举个实际例子。假设我要读从站地址为 1 的设备起始寄存器地址 0x0000读 2 个寄存器也就是 4 个字节通常对应温度和湿度各一个 16 位值。请求报文是00 01 00 00 00 06 01 03 00 00 00 02拆开看00 01是事务 ID00 00是协议 ID00 06表示后面还有 6 个字节01是从站地址03是功能码00 00是起始地址00 02是寄存器数量。响应报文大概是这样00 01 00 00 00 07 01 03 04 01 2C 02 1A00 07表示后面 7 个字节01 03是从站和功能码04是数据字节数2 个寄存器 4 字节01 2C是第一个寄存器的值0x012C 30002 1A是第二个0x021A 538。至于这 300 和 538 怎么换算成实际的温度和湿度那要看设备手册里的寄存器映射表通常是除以 10 或者除以 100。注意Modbus 的寄存器地址有个经典的坑——协议地址从 0 开始但很多设备手册写的是从 1 开始的数据地址。比如手册写温度值在 40001那对应的协议起始地址其实是 0x0000。这个 40001 里的 4 表示保持寄存器后面的 0001 才是序号。搞错这个偏移读出来的就是隔壁寄存器的数据而且往往不报错非常隐蔽。2.2 寄存器地址映射与数据解析温湿度变送器的寄存器映射各家不太一样但套路大同小异。我手上这批 POE 变送器是这样的寄存器地址协议含义数据类型换算0x0000温度int16值 / 10单位 ℃0x0001湿度uint16值 / 10单位 %RH0x0002露点int16值 / 10单位 ℃温度用int16是因为要支持零下湿度用uint16因为永远是正数。这里有个细节Python 的struct模块解析时int16要用h大端有符号uint16用H大端无符号。Modbus 协议规定所有多字节数据都是大端序也就是高字节在前这点和 x86 机器的小端序相反解析时千万别搞反。解析代码大概是这样import struct def parse_temp_humidity(data: bytes): # data 是 PDU 里的数据部分4 个字节 temp_raw, humi_raw struct.unpack(hH, data) temperature temp_raw / 10.0 humidity humi_raw / 10.0 return temperature, humidity看起来简单但实际部署时我遇到过设备返回的湿度是0xFFFF65535除以 10 变成 6553.5%RH明显是设备故障或者传感器没接好。所以解析之后一定要做合理性校验温度一般在 -40 到 85 之间湿度在 0 到 100 之间超出范围就标记为异常值别直接入库污染数据。2.3 批量读取一次请求拿多个寄存器单台设备如果只读温度和湿度一次请求读 2 个寄存器就够了。但如果设备还带露点、电池电压、信号强度等寄存器可能分散在好几个地址段。这时候有两个策略策略一合并连续地址。如果温度、湿度、露点地址是连续的 0x0000 到 0x0002那就一次请求读 3 个寄存器一个报文搞定比发三次请求快得多。策略二分段读取。如果寄存器地址跨度很大比如温度在 0x0000电池电压在 0x0100中间隔了一大片没用的地址那就分两次请求。因为 Modbus 单次读取的寄存器数量有上限通常 125 个而且读一堆没用的寄存器也是浪费带宽。我的做法是在配置阶段就把每台设备需要读的寄存器整理成段每段是一组连续地址采集时对每段发一个请求。这样既避免了无谓的读取又最大化了单次请求的效率。对于只有温湿度两台寄存器的设备就是一段一个请求。实操心得批量读取时寄存器数量宁少勿多。有些低端变送器对单次读取的寄存器数量有限制超过就返回异常码 0x03非法数据值。我一般单次不超过 10 个寄存器稳妥。3. asyncio 并发采集的完整实现3.1 单设备采集协程的编写先看最核心的单设备采集协程。这里我用asyncio.open_connection建立 TCP 连接它返回一个StreamReader和StreamWriter读写都是异步的。import asyncio import struct async def read_device(ip, port, unit_id, start_addr, reg_count, timeout2.0): try: reader, writer await asyncio.wait_for( asyncio.open_connection(ip, port), timeouttimeout ) except (asyncio.TimeoutError, OSError) as e: return {ip: ip, ok: False, error: fconnect failed: {e}} try: # 构造请求报文 transaction_id 1 pdu struct.pack(BHH, 0x03, start_addr, reg_count) mbap struct.pack(HHHB, transaction_id, 0, len(pdu) 1, unit_id) request mbap pdu writer.write(request) await asyncio.wait_for(writer.drain(), timeouttimeout) # 读 MBAP 头7 字节 header await asyncio.wait_for(reader.readexactly(7), timeouttimeout) _, _, length, _ struct.unpack(HHHB, header) # 读 PDU pdu_resp await asyncio.wait_for( reader.readexactly(length - 1), timeouttimeout ) func_code pdu_resp[0] if func_code 0x80: # 异常响应功能码最高位置 1 exception_code pdu_resp[1] return {ip: ip, ok: False, error: fmodbus exception {exception_code}} byte_count pdu_resp[1] data pdu_resp[2:2 byte_count] temp, humi struct.unpack(hH, data[:4]) return { ip: ip, ok: True, temperature: temp / 10.0, humidity: humi / 10.0 } except asyncio.TimeoutError: return {ip: ip, ok: False, error: read timeout} except Exception as e: return {ip: ip, ok: False, error: str(e)} finally: writer.close() try: await writer.wait_closed() except Exception: pass这段代码有几个关键点值得说。第一超时是分层的。连接有超时写有超时读头有超时读 PDU 也有超时。为什么要这么细因为现场网络问题五花八门有的设备 TCP 能连上但 Modbus 不响应连接超时抓不到有的响应到一半断了读 PDU 超时。分层超时能让你在日志里精确定位问题出在哪一步。第二异常响应要单独处理。Modbus 设备出错时会把功能码的最高位置 1比如 0x03 变成 0x83后面跟一个异常码。常见的异常码有0x01 非法功能、0x02 非法数据地址、0x03 非法数据值、0x04 从站设备故障。看到 0x83 就知道是读寄存器出了问题而不是网络问题。第三finally里一定要关连接。哪怕前面抛异常了writer.close()也要执行否则连接泄漏几百台设备跑几轮就把文件描述符耗光了。3.2 用 Semaphore 控制并发上限有了单设备协程最直觉的写法是给每台设备创建一个 task然后asyncio.gather全部跑起来。但这里有个大坑如果你有 500 台设备一次性创建 500 个 TCP 连接交换机和设备的连接数可能扛不住。很多低端 POE 变送器同时只支持 2 到 4 个连接你一下子砸 500 个过去大部分会被拒绝或者超时。所以必须用asyncio.Semaphore限制同时进行的连接数。我的经验值是并发数控制在 50 到 100 之间具体看交换机的性能和设备的连接能力。这个值太小采集慢太大设备扛不住。async def read_device_with_sem(sem, *args, **kwargs): async with sem: return await read_device(*args, **kwargs) async def poll_all(devices, concurrency80): sem asyncio.Semaphore(concurrency) tasks [ read_device_with_sem(sem, d[ip], d[port], d[unit_id], d[start_addr], d[reg_count]) for d in devices ] results await asyncio.gather(*tasks, return_exceptionsTrue) return resultsasyncio.gather的return_exceptionsTrue很重要。如果不加这个参数任何一个协程抛出未捕获的异常整个gather就会中断其他设备的结果全丢了。加上之后异常会作为结果返回你可以逐个检查。注意Semaphore限制的是同时持有信号量的协程数也就是同时进行的连接数。它不会限制 task 的创建数量。500 个 task 还是会全部创建出来只是大部分在async with sem那里排队等待。task 本身很轻量几百上千个完全没问题不用担心内存。3.3 连接复用短连接还是长连接上面的代码是短连接模式每次采集都新建 TCP 连接读完就关。这种方式简单、隔离性好但有个代价——每次都要经历 TCP 三次握手对于几百台设备来说握手本身的开销就不小。另一种是长连接模式每台设备维持一个常驻连接采集时直接在这个连接上发请求。这样省掉了握手开销但引入了新问题连接可能因为网络原因悄悄断开你需要心跳保活和断线重连而且几百个常驻连接对设备和服务器的资源占用都更高。我的选择是短连接为主配合连接池。具体来说用asyncio的Queue维护一个连接池采集时从池里取连接用完放回。如果连接失效了就丢弃重建。这样兼顾了效率和隔离性。不过说实话对于 200 台设备、秒级刷新的场景短连接实测下来一轮 1.5 秒已经够用了连接池的复杂度不一定值得。先跑通短连接确认性能瓶颈真的在握手上了再上连接池这是我的一贯原则。3.4 数据汇聚与入库采集回来的结果是一堆字典下一步是入库。这里我踩过一个坑不要在采集协程里直接写数据库。因为数据库写入是同步阻塞的你在协程里调用pymysql或者sqlite3会把整个事件循环卡住并发优势瞬间归零。正确做法是采集和入库分离采集协程只负责把结果放进一个asyncio.Queue另起一个消费者协程从队列里取数据、批量入库。批量入库比逐条插入快得多我一般攒够 100 条或者每隔 1 秒写一次。async def db_writer(queue, batch_size100): batch [] while True: try: item await asyncio.wait_for(queue.get(), timeout1.0) batch.append(item) if len(batch) batch_size: await asyncio.to_thread(write_batch_to_db, batch) batch.clear() except asyncio.TimeoutError: if batch: await asyncio.to_thread(write_batch_to_db, batch) batch.clear()注意这里用了asyncio.to_thread把同步的数据库写入丢到线程池里执行避免阻塞事件循环。这是 Python 3.9 之后的标准做法比手动管理ThreadPoolExecutor方便。4. 现场部署的常见问题与排查实录4.1 那些年踩过的坑坑一POE 供电不足导致设备间歇性掉线。这个坑最隐蔽。现场 24 口 POE 交换机标称总功率 370W我接了 20 台变送器每台标称 3W算下来才 60W看起来绰绰有余。但实际运行中总有几台设备随机掉线。后来查了交换机的 POE 规格才发现单口供电有上限比如 15.4W而且交换机启动瞬间的浪涌电流会拉高瞬时功率。解决办法是把设备分散到多台交换机上或者选 POE 的交换机。这个坑和软件无关但排查起来特别费劲因为现象是随机几台设备超时很容易误判成代码问题。坑二从站地址冲突。Modbus TCP 里如果设备直连不经过串口网关从站地址其实不太重要很多设备默认都是 1。但如果你的网络里有个串口服务器后面挂了好几个串口设备那从站地址就必须唯一。我有一次就是因为两台设备的从站地址都是 1导致读出来的数据张冠李戴。坑三寄存器字节序。前面说了 Modbus 是大端但有些厂家不按套路出牌把 32 位浮点数拆成两个寄存器时用了字交换word swap也就是高低字反了。这种问题只能靠对着设备手册和实际读数反复验证。我的经验是先用 Modbus Poll 这类工具手动读一次确认原始值再写代码解析别一上来就写代码猜。坑四事件循环被同步代码阻塞。有一次我在采集协程里加了个日志写入用的是普通的logging同步 handler结果并发数一上去采集速度断崖式下跌。原因是日志写文件是同步 IO把事件循环堵住了。后来换成了QueueHandler日志写入也异步化问题解决。4.2 常见问题速查表现象可能原因排查方法解决大量设备连接超时并发数过高设备连接数打满降低并发数到 20 试试调小 Semaphore 值个别设备随机超时POE 供电不稳或网线质量差换网线、换交换机端口分散供电、换 POE返回异常码 0x83 0x02寄存器地址错误对照手册核对地址偏移修正起始地址数据明显不合理字节序或换算系数错误用调试工具读原始值修正 struct 格式采集速度越来越慢连接泄漏文件描述符耗尽lsof看连接数确保 finally 关连接事件循环卡顿同步 IO 阻塞检查是否有同步库调用改用异步或 to_thread4.3 性能调优的几个实测数据我在一台 4 核 8G 的工控机上做过压测200 台设备每台读 2 个寄存器结果如下并发数单轮耗时CPU 占用失败率204.2s5%0%502.1s8%0%801.5s12%0%1501.4s18%2%2001.6s25%8%可以看到并发数从 80 往上加耗时基本不再下降反而失败率上升。80 是这个现场的最优值。这个数字不是拍脑袋定的是压测出来的。不同现场的最优值不一样取决于交换机性能、设备连接能力、网络质量一定要自己压测。实操心得压测的时候别只看总耗时一定要看失败率和P99 延迟。平均耗时好看但失败率高说明并发过头了。我一般要求失败率低于 0.5%P99 延迟低于 3 秒。5. 从采集到监控让这套系统真正跑起来5.1 定时轮询与优雅退出采集脚本不能跑一次就退出得常驻运行每隔几秒轮询一轮。这里用asyncio的循环配合信号处理实现优雅退出import signal async def main_loop(devices, interval5): stop_event asyncio.Event() def handle_signal(): stop_event.set() loop asyncio.get_running_loop() for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, handle_signal) while not stop_event.is_set(): start loop.time() results await poll_all(devices) # 处理结果... elapsed loop.time() - start wait max(0, interval - elapsed) try: await asyncio.wait_for(stop_event.wait(), timeoutwait) except asyncio.TimeoutError: pass这里有个细节等待时间要减去本轮采集耗时。如果采集花了 1.5 秒间隔是 5 秒那实际只等 3.5 秒保证轮询周期稳定。如果采集本身就超过了间隔那就立即开始下一轮不等待。add_signal_handler是 asyncio 提供的信号处理方式比传统的signal.signal更适合异步环境因为它能把信号回调安全地投递到事件循环里。5.2 异常数据的处理策略采集系统最怕的不是采不到数据而是采到了错误的数据还当成正常的入库。我的策略是三层过滤第一层协议层校验。功能码异常、字节数不对、响应长度不符直接判为失败。第二层数值范围校验。温度超出 -40 到 85湿度超出 0 到 100标记为可疑值。可疑值不入正常表进异常表方便事后分析。第三层变化率校验。如果一台设备的温度在 5 秒内从 25℃ 跳到 80℃那大概率是传感器故障或者通信错位。这种突变我会触发告警但不直接入库。这三层过滤能挡掉 95% 以上的脏数据。剩下的 5% 靠人工巡检和定期校准。5.3 后续可以扩展的方向这套框架跑通之后扩展性其实很好。比如加告警温度超过阈值通过邮件或者消息推送通知。加历史趋势数据入库后接 Grafana画温湿度曲线。加设备健康度统计每台设备的失败率失败率高的设备提前预警可能是网线老化或者设备要坏了。加动态并发根据当前失败率自动调整 Semaphore 的值网络好的时候多并发网络差的时候降下来。我个人在实际操作中的体会是采集系统 80% 的功夫在异常处理上。把正常流程跑通可能只要半天但把各种异常情况处理稳妥需要反复在现场磨。尤其是 POE 供电和网络质量这两个变量软件层面能做的有限很多时候得靠硬件选型和现场施工来保证。所以别指望一套代码打天下每个现场都要留出调试和压测的时间。最后再分享一个小技巧给每台设备加一个最后成功时间字段。如果某台设备连续 N 轮采集失败就把它暂时移出采集队列过几分钟再放回来重试。这样能避免对已经掉线的设备做无谓的轮询把并发资源留给正常设备。这个简单的机制能让整个系统在部分设备故障时依然保持高效。
返回列表