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()
问题原因及修复方案
核心问题
read_message方法的错误嵌套:线程池只能执行同步函数,但你传入了一个async函数blocking_read,且该函数内部的await self.read()无法在线程中被事件循环调度,导致整个read_message阻塞事件循环,只有串口读取超时(timeout=1)时才会让出CPU,此时WebSocket积压的消息才会被处理。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
相关产品推荐
相关产品推荐

