尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Python asyncio并发采集数百台Modbus TCP温湿度变送器实战

发布时间:2026/9/29 11:24:04

资讯中心
01
ARTICLE

Python asyncio并发采集数百台Modbus TCP温湿度变送器实战

Python asyncio并发采集数百台Modbus TCP温湿度变送器实战
1. 项目缘起与整体架构思路1.1 为什么会有这个采集需求做过机房动环监控或者仓储环境监测的朋友应该都有体会当你要监控的温湿度点位从几个变成几十个、上百个的时候事情的性质就完全变了。早期我用轮询的方式一台一台设备去读每台设备超时等个两三秒一轮下来十几分钟过去了数据的实时性根本没法保证。更别提有些POE供电的以太网温湿度变送器本身就部署在比较分散的位置网络抖动是常态串行轮询的效率低到让人抓狂。这个项目的核心目标很明确用Python的asyncio异步框架同时对数百台支持Modbus TCP/IP协议的POE以太网温湿度变送器发起并发采集把单轮采集时间从十几分钟压缩到几秒级别同时保证数据的完整性和异常处理能力。适合谁来参考如果你正在做以下事情这篇内容应该能帮到你机房、仓库、实验室的环境监控系统开发工业物联网场景下的Modbus设备批量数据采集想学习Python asyncio在真实IO密集型任务中的落地应用需要从零搭建一套可扩展的采集框架1.2 为什么选asyncio而不是多线程这里我要先解释一个关键选型问题。很多人第一反应是并发采集嘛开线程池不就行了我最初也是这么想的用concurrent.futures.ThreadPoolExecutor开50个线程去跑。实测下来有几个问题第一线程开销不可忽视。每台设备一个连接数百个线程的上下文切换成本很高GIL虽然对IO操作影响不大但线程的创建和销毁本身就有开销。第二代码复杂度高。线程间的数据汇总、异常传递、超时控制都需要额外的同步机制写起来容易出bug。第三扩展性差。当设备数量从200变成500线程模型需要重新调参而asyncio的事件循环天然支持数千个并发连接。asyncio的本质是单线程事件循环 协程调度。当某个协程在等待网络IO时事件循环会自动切换到其他就绪的协程没有线程切换的开销内存占用也极低。对于Modbus TCP这种典型的请求-响应模式每个采集任务大部分时间都在等网络返回asyncio简直是为这种场景量身定做的。1.3 整体架构设计整个采集系统的架构可以拆成四层设备管理层维护所有变送器的IP、端口、从站地址、寄存器映射表并发调度层用asyncio的Semaphore控制并发上限用TaskGroup管理任务生命周期协议通信层基于asyncio实现Modbus TCP的请求封装和响应解析数据处理层采集结果汇总、异常记录、数据入库我选择自己实现Modbus TCP协议而不是直接用pymodbus原因后面会详细说。先给一个整体的数据流设备列表加载 → 创建协程任务 → 信号量限流 → 建立TCP连接 → 发送Modbus请求 → 解析响应 → 数据校验 → 结果汇总。2. Modbus TCP协议核心细节与实操要点2.1 Modbus TCP的报文结构拆解要用asyncio手写Modbus TCP客户端必须先把协议报文吃透。Modbus TCP的报文叫MBAP头 PDU结构如下字段长度说明事务标识符2字节用于匹配请求和响应每次请求递增协议标识符2字节Modbus固定为0x0000长度2字节后续字节数单元标识符PDU单元标识符1字节从站地址TCP场景通常为0x01或0xFF功能码1字节如0x03读保持寄存器数据N字节起始地址、寄存器数量等读温湿度通常用功能码0x03读保持寄存器。假设变送器的温度值在寄存器地址0x0000湿度在0x0001那么请求PDU就是03 00 00 00 02意思是读起始地址0、连续2个寄存器。响应PDU则是03 04 [温度高字节] [温度低字节] [湿度高字节] [湿度低字节]。注意这里的字节序问题很多变送器用的是大端序但也有一些厂商用小端序或者字交换这个必须查设备手册确认。2.2 为什么不用pymodbuspymodbus确实是个成熟的库但在这个场景下我选择自己实现理由有三第一异步支持不够彻底。pymodbus的异步客户端虽然基于asyncio但内部封装层次多对于数百个并发连接每个连接都创建一个client实例资源占用比裸socket大不少。第二异常处理不够灵活。工业现场的网络环境复杂我需要精确控制每个请求的超时、重试、连接复用策略。自己写socket可以做到毫秒级的精细控制。第三调试友好。自己实现的协议层出问题时可以直接打印原始字节流排查效率高很多。用库的话出问题往往要翻源码。当然如果你对Modbus协议不熟或者项目周期紧用pymodbus也完全没问题只是性能上会打些折扣。2.3 寄存器地址的坑0还是1这是Modbus新手最容易踩的坑。Modbus协议本身规定寄存器地址从0开始但很多设备手册上写的是从1开始的地址比如40001这种。实际发送报文时需要把手册地址减1。举个例子手册上写温度值在保持寄存器40001那么实际发送的起始地址是0x0000。如果手册写40002实际发送0x0001。这个转换如果搞错了读回来的数据要么是错的要么直接报异常码0x02非法数据地址。我的做法是在设备配置表里统一用协议地址从0开始然后在备注里标注手册地址避免混淆。2.4 POE供电设备的网络特性POEPower over Ethernet供电的以太网温湿度变送器通过网线同时传输数据和电力部署非常方便一根网线搞定。但这也带来一些网络特性需要注意上电启动时间POE设备上电后需要几秒到十几秒的启动时间采集程序启动时如果设备还没就绪会大量报连接超时。建议程序启动后先等待30秒再开始采集。网络带宽POE交换机的背板带宽有限数百台设备同时通信时要注意交换机的包转发率。不过Modbus报文很小通常几十字节一般不会成为瓶颈。供电稳定性POE供电距离一般不超过100米超过后电压衰减可能导致设备工作不稳定表现为间歇性通信失败。3. asyncio并发采集的完整实现3.1 环境准备与依赖Python版本建议3.11及以上因为3.11引入了asyncio.TaskGroup任务管理更方便。如果只能用3.8那就用asyncio.gather配合return_exceptionsTrue。依赖方面我尽量用标准库只装一个uvloop来加速事件循环可选但实测能提升20%左右的吞吐pip install uvloop如果你在Windows上开发uvloop不支持直接用默认事件循环即可。生产环境建议部署在Linux上。3.2 设备配置表的设计设备信息我建议用YAML或者JSON管理方便维护。结构如下DEVICES [ { name: 机房A-01, ip: 192.168.1.101, port: 502, slave_id: 1, temp_reg: 0, # 温度寄存器协议地址 humi_reg: 1, # 湿度寄存器协议地址 scale_temp: 0.1, # 温度缩放系数 scale_humi: 0.1, # 湿度缩放系数 }, # ... 更多设备 ]这里的scale_temp和scale_humi很关键。很多变送器返回的是整数比如温度返回值235实际是23.5℃需要乘以0.1。这个系数必须查手册确认搞错了数据就全废了。3.3 核心采集协程的实现先看最核心的单个设备采集协程import asyncio import struct async def read_device(device, timeout3.0): 采集单台设备的温湿度 try: reader, writer await asyncio.wait_for( asyncio.open_connection(device[ip], device[port]), timeouttimeout ) except (asyncio.TimeoutError, OSError) as e: return {name: device[name], error: f连接失败: {e}} try: # 构造Modbus TCP请求 transaction_id 1 protocol_id 0 unit_id device[slave_id] function_code 0x03 start_addr device[temp_reg] reg_count 2 # 读温度和湿度两个寄存器 pdu struct.pack(BHH, function_code, start_addr, reg_count) length len(pdu) 1 mbap struct.pack(HHHB, transaction_id, protocol_id, length, unit_id) request mbap pdu writer.write(request) await writer.drain() # 读取响应MBAP头7字节 PDU response await asyncio.wait_for(reader.read(256), timeouttimeout) if len(response) 9: return {name: device[name], error: 响应过短} # 解析响应 resp_func response[7] if resp_func 0x80: # 异常响应 exception_code response[8] return {name: device[name], error: fModbus异常码: {exception_code}} byte_count response[8] data response[9:9byte_count] if len(data) 4: return {name: device[name], error: 数据长度不足} temp_raw struct.unpack(h, data[0:2])[0] humi_raw struct.unpack(h, data[2:4])[0] return { name: device[name], temperature: temp_raw * device[scale_temp], humidity: humi_raw * device[scale_humi], } except asyncio.TimeoutError: return {name: device[name], error: 读取超时} except Exception as e: return {name: device[name], error: f解析异常: {e}} finally: writer.close() try: await writer.wait_closed() except Exception: pass这段代码有几个关键点需要说明事务标识符的处理上面简化用了固定值1实际生产环境应该用递增计数器避免响应错配。不过因为每个协程用的是独立连接请求和响应是一一对应的固定值也不会出问题。超时控制连接超时和读取超时都设了3秒。这个值需要根据网络状况调整局域网内1秒足够跨网段建议3-5秒。异常响应判断功能码最高位为1表示异常响应比如请求的功能码是0x03异常响应就是0x83。异常码常见的有0x01非法功能、0x02非法数据地址、0x03非法数据值。3.4 并发调度与限流有了单设备采集协程接下来就是并发调度。这里的关键是信号量限流不能一次性发起数百个连接否则会打爆本机文件描述符或者POE交换机。async def collect_all(devices, max_concurrent100): 并发采集所有设备 semaphore asyncio.Semaphore(max_concurrent) async def limited_read(device): async with semaphore: return await read_device(device) tasks [limited_read(d) for d in devices] results await asyncio.gather(*tasks, return_exceptionsTrue) # 处理异常结果 final [] for r in results: if isinstance(r, Exception): final.append({error: str(r)}) else: final.append(r) return finalmax_concurrent设多少合适我的经验值是50到100。设太小浪费并发能力设太大可能导致本机端口耗尽或者交换机拥塞。可以先从50开始压测逐步往上加观察采集成功率和耗时。如果你用的是Python 3.11可以用TaskGroup替代gather异常处理更优雅async def collect_all_v2(devices, max_concurrent100): semaphore asyncio.Semaphore(max_concurrent) results [] async def limited_read(device): async with semaphore: result await read_device(device) results.append(result) async with asyncio.TaskGroup() as tg: for d in devices: tg.create_task(limited_read(d)) return results3.5 连接复用与性能优化上面的实现每次采集都新建TCP连接对于需要高频采集的场景比如每10秒一轮连接建立的开销占比不小。优化方案是连接池为每台设备维护一个长连接采集时复用。但长连接有个问题POE设备可能因为网络抖动断开复用时会报错。我的做法是连接失败后自动降级为新建连接并在下一轮采集时重建长连接。class DeviceConnection: def __init__(self, device): self.device device self.reader None self.writer None self.lock asyncio.Lock() async def ensure_connected(self): if self.writer is None or self.writer.is_closing(): self.reader, self.writer await asyncio.wait_for( asyncio.open_connection(self.device[ip], self.device[port]), timeout3.0 ) async def close(self): if self.writer: self.writer.close() try: await self.writer.wait_closed() except Exception: pass self.writer None实测下来连接复用能把单轮采集时间再压缩30%左右对于数百台设备的场景这个提升很可观。4. 常见问题与排查技巧实录4.1 采集失败问题速查表现象可能原因排查方法解决方案连接超时设备未上电/IP错误/网络不通ping设备IP检查POE供电和网线连接被拒绝端口错误/设备未监听502telnet IP 502确认设备Modbus TCP端口响应超时网络延迟大/设备繁忙抓包看请求是否到达增大超时时间异常码0x02寄存器地址错误核对手册地址地址减1重试数据明显异常字节序错误/缩放系数错误对比设备显示屏调整字节序或系数间歇性失败POE供电不稳/网络抖动观察失败设备分布检查POE交换机功率4.2 字节序问题的排查字节序是Modbus采集里最隐蔽的坑。同一个温度值23.5℃不同厂商的返回可能完全不同大端序00 EB235小端序EB 00字交换如果读两个寄存器温度和湿度的顺序可能颠倒排查方法用设备自带的显示屏或者厂商调试软件读一次记录原始值然后对比你的程序读到的原始字节。如果对不上就调整struct.unpack的格式符h是大端有符号h是小端有符号。4.3 并发数调优的实操经验我做过一组压测200台设备不同并发数下的表现并发数单轮耗时成功率备注2012.3秒99.5%太保守505.1秒99.2%推荐起点1003.2秒98.7%性能最佳2003.0秒95.3%开始丢包5004.8秒88.1%交换机拥塞结论并发数不是越大越好超过交换机处理能力后成功率下降反而拖慢整体。建议从50开始逐步加到成功率开始明显下降为止。4.4 几个容易忽略的细节第一文件描述符限制。Linux默认单进程最多1024个文件描述符每个TCP连接占一个。如果并发数设到500以上需要调整ulimit -n。第二DNS解析。如果设备配置里用主机名而不是IP每次连接都要DNS解析会拖慢速度。建议全部用IP。第三日志写入。数百个协程同时写日志文件会有IO竞争建议用asyncio.Queue做异步日志或者直接写内存后批量落盘。第四采集周期与超时的关系。如果采集周期是10秒单设备超时设3秒那么最坏情况下一个设备要等3秒。要保证单轮采集时间小于周期否则任务会堆积。5. 数据校验与异常处理策略5.1 数据合理性校验采集回来的数据不能直接入库必须做合理性校验。温湿度变送器的合理范围一般是温度-40℃ 到 85℃工业级湿度0% 到 100%超出这个范围的数据要么是传感器故障要么是通信错误导致的解析异常。我的做法是标记为异常数据记录原始字节但不入库同时触发告警。def validate(temp, humi): if not (-40 temp 85): return False, f温度超范围: {temp} if not (0 humi 100): return False, f湿度超范围: {humi} return True, None5.2 重试策略的设计对于偶发的通信失败重试是必要的。但重试不能无脑重试要有策略连接失败立即重试1次间隔0.5秒读取超时重试1次间隔1秒异常码响应不重试直接记录说明是配置问题重试也没用重试次数不宜过多2次足够。因为如果设备真的挂了重试只会拖慢整体采集速度。5.3 采集结果的汇总与告警每轮采集结束后我会统计以下指标成功数、失败数、成功率平均耗时、最大耗时失败设备列表及原因如果成功率低于95%或者连续3轮某台设备失败就触发告警。告警方式可以是邮件、钉钉机器人或者写入告警表。6. 部署与长期运行的经验6.1 用systemd管理采集服务生产环境建议用systemd托管配置如下[Unit] DescriptionModbus TCP Collector Afternetwork.target [Service] Typesimple Usercollector WorkingDirectory/opt/modbus-collector ExecStart/opt/modbus-collector/venv/bin/python main.py Restartalways RestartSec10 LimitNOFILE65535 [Install] WantedBymulti-user.targetLimitNOFILE65535这行很关键解决文件描述符限制问题。Restartalways保证程序崩溃后自动重启。6.2 内存泄漏的预防asyncio程序长期运行最容易出的问题是未关闭的连接和未取消的任务。我的经验是每个连接用完必须close()并在finally块里执行用asyncio.all_tasks()定期检查是否有僵尸任务采集主循环加一个总超时防止某轮采集卡死6.3 采集频率的权衡采集频率越高数据越实时但对设备和网络的负担越大。温湿度变化本身比较缓慢30秒到60秒一轮完全够用。如果是机房精密空调监控可以缩短到10秒。不建议低于5秒因为很多变送器的采样周期本身就有几秒采太快没意义。7. 性能实测与扩展方向7.1 实测数据在200台设备的真实环境中我的最终配置是并发100、超时3秒、连接复用实测单轮采集耗时稳定在2.8到3.5秒成功率99%以上。相比最初的串行轮询约15分钟提升了近300倍。CPU占用方面采集进程常驻内存约80MBCPU峰值不到15%非常轻量。7.2 后续可以扩展的方向这套框架的扩展性很好几个可以继续做的方向数据入库接入InfluxDB或TDengine做时序数据存储和可视化动态设备发现扫描网段自动发现Modbus TCP设备协议扩展同一套框架支持Modbus RTU over TCP兼容更多设备边缘计算在采集端做数据预处理只上传异常和统计值我个人在实际操作中的体会是asyncio这套东西入门曲线有点陡但一旦理解了事件循环和协程调度的本质写起来比多线程舒服太多。尤其是异常处理每个协程独立try-except互不影响这在工业采集场景下太重要了。踩过的坑主要集中在字节序和寄存器地址上这两个问题解决了剩下的就是调参和优化的事。
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。