You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python实现WebSocket转串口桥接程序出现阻塞问题

问题:WebSocket转串口桥接程序中WebSocket消息无法及时转发到串口

我用Python的websockets和pyserial库编写WebSocket转串口的桥接程序,目标是:

  • 将WebSocket收到的消息转发至串口设备
  • 持续读取串口消息(以换行符为结束标志)并转发至WebSocket

当前遇到的问题:

  • 串口读取已放在线程池执行,能正常将串口消息实时转发到WebSocket
  • WebSocket收到的消息无法及时转发到串口,只有当串口收到新消息后,积压的WebSocket消息才会批量发送,无法做到实时转发

代码实现

SerialConnection类

from serial import Serial
import asyncio
from .serial_emulator import SerialEmulator
from concurrent.futures import ThreadPoolExecutor


class SerialConnection:
    def __init__(self, port, baudrate=115200, timeout=1):
        if port == "test":
            self.connection = SerialEmulator()
        else:
            self.connection = Serial(port, baudrate=baudrate, timeout=timeout)

        self.reader = asyncio.StreamReader()
        self.executor = ThreadPoolExecutor()

    def write(self, data):
        self.connection.write(data)

    async def read(self, size: int = 1):
        return self.connection.read(size)

    def close(self):
        self.connection.close()

    @property
    def port(self):
        if self.connection is None:
            return None
        return self.connection.port

    async def read_message(self):
        if self.connection is None:
            return

        async def blocking_read():
            message = b""
            while not message.endswith(b"\n"):
                data = await self.read()
                message += data
            return message.decode().rstrip("\r\n")

        loop = asyncio.get_running_loop()
        message = await loop.run_in_executor(self.executor, blocking_read)
        return await message

WebSocketServer类

import asyncio
import websockets


class WebSocketServer:
    def __init__(self, port, handler=None):
        self.port = port
        self.server = None
        self.handler = handler or self.handler

    async def handler(self, websocket, path):
        async for message in websocket:
            print(f"websocket message: {message}, path: {path}")

    async def start(self):
        self.server = await websockets.serve(self.handler, "localhost", self.port)
        await self.server.wait_closed()

    async def stop(self):
        if self.server:
            self.server.close()
            await self.server.wait_closed()

    def run_forever(self):
        asyncio.get_event_loop().run_until_complete(self.start())
        asyncio.get_event_loop().run_forever()

main函数

import argparse
import asyncio

from websocket_serial_bridge import WebSocketServer, SerialConnection


def main():
    parser = argparse.ArgumentParser()

    parser.add_argument("-p", "--websocket-port", type=int, default=8765)

    # 列出所有可用串口后退出
    parser.add_argument("-l", "--list", "--list-serial", action="store_true")

    # 指定要连接的串口列表
    parser.add_argument("-s", "--serial-ports", nargs="+", type=str, default=[])

    # 测试模式(模拟串口)
    parser.add_argument("-t", "--test", action="store_true", help="Run in test mode (emulate serial ports).")

    args = parser.parse_args()

    def list_serial_ports():
        import serial.tools.list_ports
        ports = serial.tools.list_ports.comports()
        if len(ports) == 0:
            print("未找到串口")
            return
        for port in ports:
            print(f" - {port.device} - {port.description}")

    if args.list:
        print("列出可用串口:")
        list_serial_ports()
        return

    websocket_port = args.websocket_port

    serial_ports = args.serial_ports
    if args.test:
        if len(serial_ports) > 0:
            print("警告:已指定测试模式,忽略串口参数")
        serial_ports = ["test"]

    # 未指定串口时列出所有可用串口
    if len(serial_ports) == 0:
        print("未指定串口,列出所有可用串口:")
        list_serial_ports()
        return

    serial_connections = []
    for serial_port in serial_ports:
        serial_connections.append(SerialConnection(serial_port))

    async def serial_to_websocket(websocket, serial_connection):
        while True:
            serial_msg = await serial_connection.read_message()
            message = f"PORT:{serial_connection.port}|{serial_msg}"
            print(message)
            await websocket.send(serial_msg)

    async def websocket_to_serial(websocket, serial_connection):
        async for websocket_msg in websocket:
            print(f"从WebSocket收到消息:{websocket_msg}")
            serial_connection.write(websocket_msg.encode() + b"\n")

    async def handler(websocket, path):
        serial_connection = serial_connections[0]

        await asyncio.gather(
            asyncio.create_task(websocket_to_serial(websocket, serial_connection)),
            asyncio.create_task(serial_to_websocket(websocket, serial_connection)),
        )

    server = WebSocketServer(websocket_port, handler=handler)

    print(f"WebSocket服务启动:ws://localhost:{websocket_port}")

    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)

    loop.create_task(server.start())
    loop.run_forever()


if __name__ == "__main__":
    main()

问题原因及修复方案

核心问题

  1. read_message方法的错误嵌套:线程池只能执行同步函数,但你传入了一个async函数blocking_read,且该函数内部的await self.read()无法在线程中被事件循环调度,导致整个read_message阻塞事件循环,只有串口读取超时(timeout=1)时才会让出CPU,此时WebSocket积压的消息才会被处理。
  2. read方法是伪异步:self.connection.read(size)是同步阻塞调用,但你直接包装成async方法却未放到线程池执行,会直接阻塞事件循环。

修复步骤

1. 修复SerialConnection类的读取逻辑

将串口的同步操作完全隔离到线程池,避免异步嵌套:

from serial import Serial
import asyncio
from .serial_emulator import SerialEmulator
from concurrent.futures import ThreadPoolExecutor


class SerialConnection:
    def __init__(self, port, baudrate=115200, timeout=1):
        if port == "test":
            self.connection = SerialEmulator()
        else:
            self.connection = Serial(port, baudrate=baudrate, timeout=timeout)

        self.executor = ThreadPoolExecutor(max_workers=1)  # 串口读取单线程足够

    def write(self, data):
        self.connection.write(data)

    def _read_sync(self, size: int = 1):
        """同步读取串口数据"""
        return self.connection.read(size)

    async def read(self, size: int = 1):
        """异步封装同步读取"""
        loop = asyncio.get_running_loop()
        return await loop.run_in_executor(self.executor, self._read_sync, size)

    def close(self):
        self.connection.close()
        self.executor.shutdown()

    @property
    def port(self):
        return self.connection.port if self.connection else None

    async def read_message(self):
        if not self.connection:
            return

        def blocking_read():
            """完全同步的读取逻辑,放到线程池执行"""
            message = b""
            while not message.endswith(b"\n"):
                data = self.connection.read(1)
                if not data:
                    # 超时后返回已读取内容(如果有)
                    break
                message += data
            return message.decode().rstrip("\r\n") if message else None

        loop = asyncio.get_running_loop()
        return await loop.run_in_executor(self.executor, blocking_read)

2. 优化WebSocket消息的写入并发(可选)

如果串口写入存在阻塞风险(如缓冲区满),可以把写入操作也放到线程池:

# 在SerialConnection类中添加异步写入方法
async def write_async(self, data):
    loop = asyncio.get_running_loop()
    await loop.run_in_executor(self.executor, self.connection.write, data)

修改websocket_to_serial函数:

async def websocket_to_serial(websocket, serial_connection):
    async for websocket_msg in websocket:
        print(f"从WebSocket收到消息:{websocket_msg}")
        await serial_connection.write_async(websocket_msg.encode() + b"\n")

3. 简化事件循环启动

替换main函数中的循环启动代码,使用标准的asyncio.run:

# 替换原来的循环启动代码
asyncio.run(server.start())

修复后效果

  • 串口读取在独立线程中执行,不再阻塞事件循环
  • WebSocket消息收到后立即转发到串口,无延迟
  • 两个方向的数据流完全独立并发,互不干扰

内容的提问来源于stack exchange,提问作者Trantidon

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 05:30:12