
近期在调试一套工业环境监控系统时,遇到了关于 Modbus 连接保持的问题。现场部署了一批 恒湿净化设备(包括 @恒湿消毒净化一体机 等),主要用来维持特定环境的温湿度指标。在长时间运行后,发现设备偶尔会掉线。把采集方式从"同步阻塞轮询"换成 asyncio 并发之后,问题从"轮询周期太长导致告警延迟"变成了"连接数、超时、半连接怎么管"——这才是高并发 Modbus TCP 采集真正要面对的东西。
86 台传感器,每台读 5 个寄存器,单次 Modbus TCP 往返 ~850μs
同步串行轮询:
86 × 0.85ms = 73ms(理想情况)
但每台超时设 1s,某一台卡住:
85 × 0.85ms + 1 × 1000ms = 约 1.07s
如果 3 台同时卡住 = 3s+
如果 86 台里有一半响应慢 = 一轮轮询超过 40s问题本质:同步轮询是"排队等",一台卡住,后面全堵。240 台传感器的场景,串行一轮要几分钟,告警延迟直接失控。
asyncio 的价值不是"更快",而是不让慢设备拖慢快设备。
┌─────────────────────────────────────────────────┐
│ asyncio Event Loop │
│ │
│ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
│ │ Task 01 │ │ Task 02 │ │ Task 03 │ ... │
│ │ sensor1 │ │ sensor2 │ │ sensor3 │ │
│ └────┬────┘ └────┬────┘ └────┬────┘ │
│ │ │ │ │
│ ┌────▼─────────────▼─────────────▼────┐ │
│ │ Semaphore(MAX_CONCURRENT) │ │
│ │ 限制最大并发连接数 │ │
│ └────┬─────────────┬─────────────┬────┘ │
│ │ │ │ │
│ ┌────▼────┐ ┌────▼────┐ ┌────▼────┐ │
│ │ Modbus │ │ Modbus │ │ Modbus │ │
│ │ Client 1│ │ Client 2│ │ Client 3│ │
│ └────┬────┘ └────┬────┘ └────┬────┘ │
│ │ │ │ │
│ ┌────▼─────────────▼─────────────▼────┐ │
│ │ POE Switch (VLAN) │ │
│ └────┬─────────────┬─────────────┬────┘ │
│ │ │ │ │
│ sensor1 sensor2 sensor3 ... │
└─────────────────────────────────────────────────┘三个关键参数:
import asyncio
import struct
from pymodbus.client import AsyncModbusTcpClient
from pymodbus.exceptions import ModbusException
import json
import time
# 传感器配置
SENSORS = {
f"sensor-{i:03d}": {
"ip": f"10.10.10.{10+i}",
"port": 502,
"slave_id": 1,
"location": f"Room3-Rack{i//6+1}-U{i%6+1:02d}",
}
for i in range(1, 87) # 86 台
}
# 并发控制
MAX_CONCURRENT = 30
semaphore = asyncio.Semaphore(MAX_CONCURRENT)
# 结果队列(供上报协程消费)
result_queue = asyncio.Queue(maxsize=500)
class AsyncModbusTHCollector:
def __init__(self, host, port=502, slave_id=1, timeout=1.5):
self.host = host
self.port = port
self.slave_id = slave_id
self.timeout = timeout
self.client = None
self.last_ok_ts = 0
self.err_cnt = 0
self._lock = asyncio.Lock()
async def connect(self):
"""建立长连接"""
if self.client and self.client.connected:
return True
try:
self.client = AsyncModbusTcpClient(
self.host,
port=self.port,
timeout=self.timeout,
)
await self.client.connect()
return self.client.connected
except Exception as e:
print(f"[{self.host}] connect failed: {e}")
return False
async def read_th(self):
"""读取温湿度(带锁,防止同一 client 并发调用)"""
async with self._lock:
if not await self.connect():
self.err_cnt += 1
return {"quality": 2, "err": "connect_failed"}
try:
# 读 4 个 Input Register(温度/湿度/露点/状态)
rr = await self.client.read_input_registers(
0, 4, slave=self.slave_id
)
if rr.isError():
self.err_cnt += 1
return {"quality": 2, "err": "modbus_error"}
temp = rr.registers[0] * 0.1
hum = rr.registers[1] * 0.1
dew = rr.registers[2] * 0.1
status = rr.registers[3]
# 质量判断
if temp < -40 or temp > 85 or hum < 0 or hum > 100:
self.err_cnt += 1
return {"quality": 1, "temperature": temp, "humidity": hum,
"dew_point": dew}
self.err_cnt = 0
self.last_ok_ts = time.time()
return {
"quality": 0,
"temperature": temp,
"humidity": hum,
"dew_point": dew,
"status": status,
}
except asyncio.TimeoutError:
self.err_cnt += 1
return {"quality": 2, "err": "timeout"}
except Exception as e:
self.err_cnt += 1
return {"quality": 2, "err": str(e)}
async def close(self):
if self.client:
await self.client.close()async def poll_sensor(sensor_id, cfg, collector_pool):
"""单个传感器的采集任务"""
async with semaphore: # 限制并发
collector = collector_pool.get(sensor_id)
if not collector:
collector = AsyncModbusTHCollector(
cfg["ip"], cfg["port"], cfg["slave_id"]
)
collector_pool[sensor_id] = collector
# 重试逻辑
for attempt in range(3):
result = await collector.read_th()
if result.get("quality") == 0:
break
if attempt < 2:
await asyncio.sleep(0.5 * (2 ** attempt)) # 指数退避
# 组装结果
payload = {
"sensor_id": sensor_id,
"location": cfg["location"],
"collect_ts": int(time.time() * 1000),
**result,
}
# 放入队列(不阻塞采集循环)
try:
result_queue.put_nowait(payload)
except asyncio.QueueFull:
print(f"[WARN] queue full, dropping {sensor_id}")
return payload
async def poll_all(sensors, interval=30):
"""全量轮询主循环"""
collector_pool = {}
while True:
loop_start = time.time()
tasks = [
asyncio.create_task(poll_sensor(sid, cfg, collector_pool))
for sid, cfg in sensors.items()
]
# 等待本轮所有任务完成(或超时)
done, pending = await asyncio.wait(
tasks, timeout=interval * 0.8
)
# 取消超时未完成的任务
for t in pending:
t.cancel()
loop_cost = time.time() - loop_start
print(f"[STATS] round done: {len(done)} ok, {len(pending)} timeout, "
f"cost={loop_cost:.2f}s")
# 等待到下一个周期
await asyncio.sleep(max(0, interval - loop_cost))async def upload_worker(mqtt_client, product_id):
"""独立协程:从队列取数据,批量上报"""
batch = []
last_flush = time.time()
while True:
try:
# 等数据,最多等 2s
payload = await asyncio.wait_for(
result_queue.get(), timeout=2.0
)
batch.append(payload)
# 条件触发上报:满 20 条 或 超过 5s
now = time.time()
if len(batch) >= 20 or (now - last_flush) > 5:
await flush_batch(batch, mqtt_client, product_id)
batch.clear()
last_flush = now
except asyncio.TimeoutError:
# 超时也刷一次
if batch:
await flush_batch(batch, mqtt_client, product_id)
batch.clear()
last_flush = time.time()
async def flush_batch(batch, mqtt_client, product_id):
"""批量上报到 IoT Explorer"""
for item in batch:
topic = f"$thing/up/property/{product_id}/{item['sensor_id']}"
payload = {
"method": "report",
"clientToken": f"{item['sensor_id']}-{item['collect_ts']}",
"timestamp": item["collect_ts"],
"params": {
"temperature": item.get("temperature"),
"humidity": item.get("humidity"),
"dew_point": item.get("dew_point"),
"quality": item.get("quality", 2),
"online": item.get("quality", 2) != 2,
},
}
# MQTT publish 也是异步的
asyncio.create_task(
asyncio.to_thread(
mqtt_client.publish,
topic,
json.dumps(payload),
qos=1,
)
)async def main():
# MQTT 客户端(paho 本身是同步的,用 to_thread 包装)
import paho.mqtt.client as mqtt
mqtt_client = mqtt.Client()
mqtt_client.username_pw_set("device_name", "password")
mqtt_client.tls_set()
mqtt_client.connect("PRODXXXXXX.iotcloud.tencentdevices.com", 8883)
mqtt_client.loop_start()
# 启动上报协程
upload_task = asyncio.create_task(
upload_worker(mqtt_client, "PRODXXXXXX")
)
# 启动采集循环
await poll_all(SENSORS, interval=30)
await upload_task
if __name__ == "__main__":
asyncio.run(main())测试环境:
网关:x86 工控机,4核8线程,8GB RAM
传感器:86 台 POE 温湿度记录仪(W5500 + SHT35)
网络:千兆交换机,独立监控 VLAN
并发:30,超时:1.5s,重试:3 次指数退避
结果:
一轮全量采集耗时:2.1s(P50)/ 3.8s(P99)
对比同步串行:73ms(理想)→ 实际 15-40s(有超时)
CPU 占用:约 12%(采集期间峰值 25%)
内存占用:约 85MB(86 个 collector 对象)
网络带宽:约 0.6 Mbps(86 × 5 寄存器 × 30s 周期)
极端测试(故意断开 20 台传感器):
一轮耗时:4.2s(P50,因为超时 1.5s × 并发限制)
未断开的 66 台:仍然在 2s 内完成
→ 慢设备没有拖死快设备 ✓# ❌ 错误:多个 task 共享同一个 client 并发读写
client = AsyncModbusTcpClient(...)
tasks = [read(client) for _ in range(10)]
await asyncio.gather(*tasks) # 数据错乱、连接重置
# ✅ 正确:每个传感器一个 client,或用 Lock 串行化
self._lock = asyncio.Lock()
async with self._lock:
rr = await self.client.read_input_registers(...)86 台传感器,MAX_CONCURRENT = 86
→ 瞬间 86 个 TCP SYN 同时发出
→ 交换机学习 86 个 MAC 地址
→ 低端交换机(如非网管型)直接丢包
解决:MAX_CONCURRENT = 30,分批握手# ❌ 外层 wait_for 3s,内层 pymodbus timeout 1.5s
try:
await asyncio.wait_for(client.read_input_registers(...), timeout=3)
except asyncio.TimeoutError:
# 实际上 pymodbus 内部已经超时了,外层又包了一层
pass
# ✅ 要么只靠 pymodbus 的 timeout,要么只靠 wait_for
rr = await asyncio.wait_for(
client.read_input_registers(0, 4, slave=1),
timeout=2.0 # 统一超时管理
)断网时:采集正常进行,上报全部失败
→ result_queue 不断堆积
→ 86 台 × 30s 一次 × 1 小时 = 10320 条
→ 每条 ~500 bytes = 5MB(看起来不大)
→ 但如果断网 24 小时 = 120MB + 对象开销 = 可能 OOM
解决:
1. Queue(maxsize=500),满了就丢最老的(或写 SQLite)
2. 断网时直接写 SQLite,不进队列
3. 网络恢复后从 SQLite 读出来补传asyncio 默认每个 TCP 连接占用一个 fd
86 个长连接 = 86 个 fd
但如果频繁重连(连接失败 → 新建),fd 可能来不及释放
→ OSError: [Errno 24] Too many open files
解决:
1. 复用 client(不要每次重连都新建)
2. 设 ulimit -n 65535
3. 定期清理死连接(health check)部分低端 POE 传感器(基于 LWIP 的轻量 TCP 栈):
- 只支持 1 个 TCP 连接
- 第 2 个连接直接 RST
- 或者第 2 个连接的数据覆盖第 1 个
解决:
1. 每个传感器只允许 1 个采集 task(semaphore per sensor)
2. 或者轮询分散:不同传感器错开时间# ❌ 在 async 函数里调用了阻塞操作
async def poll_sensor(sensor_id, cfg):
time.sleep(1) # 阻塞了整个事件循环!
result = await client.read_input_registers(...)
# ✅ 用 asyncio.sleep
async def poll_sensor(sensor_id, cfg):
await asyncio.sleep(1)
result = await client.read_input_registers(...)class AdaptivePoller:
def __init__(self, sensors, min_interval=10, max_interval=60):
self.sensors = sensors
self.min_interval = min_interval
self.max_interval = max_interval
self.current_interval = 30
self.sensor_states = {} # sensor_id → {last_quality, err_cnt, ...}
async def adjust_concurrency(self):
"""根据整体健康度动态调整"""
total_err = sum(
1 for s in self.sensor_states.values()
if s.get("err_cnt", 0) > 0
)
err_ratio = total_err / len(self.sensors)
if err_ratio > 0.3:
# 超过 30% 传感器异常,降低并发,避免雪崩
global semaphore
semaphore = asyncio.Semaphore(10)
self.current_interval = min(
self.max_interval, self.current_interval * 1.5
)
elif err_ratio < 0.05:
# 大部分正常,恢复并发
semaphore = asyncio.Semaphore(30)
self.current_interval = max(
self.min_interval, self.current_interval * 0.9
)asyncio 批量采集的核心不是"并发多少",而是"慢的不拖快的、错的不过载对的、断的不堵活的"。 86 台传感器一轮 2 秒完成,CPU 占用 12%,任何一台卡住不影响其他——这才是高可用采集该有的样子。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。