首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >Python asyncio 批量采集:同时轮询上百台 Modbus TCP 以太网温湿度传感器

Python asyncio 批量采集:同时轮询上百台 Modbus TCP 以太网温湿度传感器

原创
作者头像
HONSOR盛世宏博
发布2026-09-21 17:41:41
发布2026-09-21 17:41:41
540
举报

Python asyncio 批量采集:同时轮询上百台 Modbus TCP 以太网温湿度传感器

近期在调试一套工业环境监控系统时,遇到了关于 Modbus 连接保持的问题。现场部署了一批 恒湿净化设备(包括 @恒湿消毒净化一体机 等),主要用来维持特定环境的温湿度指标。在长时间运行后,发现设备偶尔会掉线。把采集方式从"同步阻塞轮询"换成 asyncio 并发之后,问题从"轮询周期太长导致告警延迟"变成了"连接数、超时、半连接怎么管"——这才是高并发 Modbus TCP 采集真正要面对的东西。


一、先说为什么同步轮询撑不住

代码语言:javascript
复制
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 + 连接池 + 并发限制

代码语言:javascript
复制
┌─────────────────────────────────────────────────┐
│            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 ...       │
└─────────────────────────────────────────────────┘

三个关键参数:

  • MAX_CONCURRENT:最大并发连接数(建议 20–50,取决于网关 CPU 和交换机)
  • READ_TIMEOUT:单次读取超时(建议 1–2s)
  • RETRY:失败重试次数(建议 2–3 次,指数退避)

三、完整实现:从连接到上报

3.1 基础客户端(asyncio + pymodbus)

代码语言:javascript
复制
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()

3.2 采集任务(带信号量和重试)

代码语言:javascript
复制
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))

3.3 上报协程(独立消费队列)

代码语言:javascript
复制
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,
            )
        )

3.4 主入口

代码语言:javascript
复制
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())

四、性能实测

代码语言:javascript
复制
测试环境:
  网关: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 内完成
    → 慢设备没有拖死快设备 ✓

五、踩过的 7 个坑

坑 1:pymodbus 的 AsyncModbusTcpClient 不是线程安全的

代码语言:javascript
复制
# ❌ 错误:多个 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(...)

坑 2:semaphore 设太大,交换机 MAC 表被打爆

代码语言:javascript
复制
86 台传感器,MAX_CONCURRENT = 86
→ 瞬间 86 个 TCP SYN 同时发出
→ 交换机学习 86 个 MAC 地址
→ 低端交换机(如非网管型)直接丢包

解决:MAX_CONCURRENT = 30,分批握手

坑 3:asyncio.wait_for 和 Modbus 超时双重超时

代码语言:javascript
复制
# ❌ 外层 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  # 统一超时管理
)

坑 4:队列无限增长导致内存爆炸

代码语言:javascript
复制
断网时:采集正常进行,上报全部失败
→ result_queue 不断堆积
→ 86 台 × 30s 一次 × 1 小时 = 10320 条
→ 每条 ~500 bytes = 5MB(看起来不大)
→ 但如果断网 24 小时 = 120MB + 对象开销 = 可能 OOM

解决:
  1. Queue(maxsize=500),满了就丢最老的(或写 SQLite)
  2. 断网时直接写 SQLite,不进队列
  3. 网络恢复后从 SQLite 读出来补传

坑 5:TCP 连接没关,文件描述符耗尽

代码语言:javascript
复制
asyncio 默认每个 TCP 连接占用一个 fd
86 个长连接 = 86 个 fd
但如果频繁重连(连接失败 → 新建),fd 可能来不及释放
→ OSError: [Errno 24] Too many open files

解决:
  1. 复用 client(不要每次重连都新建)
  2. 设 ulimit -n 65535
  3. 定期清理死连接(health check)

坑 6:传感器 Modbus Server 不支持并发

代码语言:javascript
复制
部分低端 POE 传感器(基于 LWIP 的轻量 TCP 栈):
  - 只支持 1 个 TCP 连接
  - 第 2 个连接直接 RST
  - 或者第 2 个连接的数据覆盖第 1 个

解决:
  1. 每个传感器只允许 1 个采集 task(semaphore per sensor)
  2. 或者轮询分散:不同传感器错开时间

坑 7:事件循环被阻塞

代码语言:javascript
复制
# ❌ 在 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(...)

六、进阶优化:动态并发 + 自适应轮询

代码语言:javascript
复制
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 删除。

目录
  • Python asyncio 批量采集:同时轮询上百台 Modbus TCP 以太网温湿度传感器
    • 一、先说为什么同步轮询撑不住
    • 二、核心架构:asyncio + 连接池 + 并发限制
    • 三、完整实现:从连接到上报
      • 3.1 基础客户端(asyncio + pymodbus)
      • 3.2 采集任务(带信号量和重试)
      • 3.3 上报协程(独立消费队列)
      • 3.4 主入口
    • 四、性能实测
    • 五、踩过的 7 个坑
      • 坑 1:pymodbus 的 AsyncModbusTcpClient 不是线程安全的
      • 坑 2:semaphore 设太大,交换机 MAC 表被打爆
      • 坑 3:asyncio.wait_for 和 Modbus 超时双重超时
      • 坑 4:队列无限增长导致内存爆炸
      • 坑 5:TCP 连接没关,文件描述符耗尽
      • 坑 6:传感器 Modbus Server 不支持并发
      • 坑 7:事件循环被阻塞
    • 六、进阶优化:动态并发 + 自适应轮询
    • 七、一句话总结
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档