串口通信高层抽象实现:多线程并发查询冲突解决问询
串口设备多线程并发查询的通用抽象实现
问题背景
在具备串口通信的硬件设备上,需要持续每秒查询一次“foo”(例如记录设备温度),同时另一线程可能在随机时机查询“bar”。常规方案是中断“foo”查询、执行“bar”查询后重启“foo”,但需要更通用的解决方案。
能否实现一个高层抽象serial_device_abstraction,让开发者只需编写以下代码,由抽象层处理底层并发问题?
def thread1(): while True: serial_device_abstraction.get("foo") time.sleep(1) def thread2(): time.sleep(random.random()) serial_device_abstraction.get("bar") threading.Thread(target=thread1).start() threading.Thread(target=thread2).start()
注:直接使用如下
get方法会因多线程同时读写串口导致数据损坏:def get(query): serial_port.write(query) data = self.serial_port.read(8)
解决方案
完全可以实现这个高层抽象,核心是通过线程安全机制保护串口资源,避免多线程同时读写。以下提供两种实用实现方案:
方案1:基础线程安全版本(锁机制)
用threading.Lock包裹串口的读写逻辑,确保同一时间只有一个线程能访问串口。这种实现简单直接,适合大多数场景,且不会中断正在执行的查询(比中断重启方案更稳定)。
import threading import serial import time import random class SerialDeviceAbstraction: def __init__(self, port, baudrate=9600): self.serial_port = serial.Serial(port, baudrate) self._serial_lock = threading.Lock() # 线程锁保护串口操作 def get(self, query): # 自动加锁/释放锁,确保同一时间仅一个线程执行串口读写 with self._serial_lock: self.serial_port.write(query.encode()) # 假设query为字符串,转字节发送 data = self.serial_port.read(8) return data.decode() # 根据实际需求处理返回数据 # 实例化抽象层(替换为你的串口路径) serial_device_abstraction = SerialDeviceAbstraction("/dev/ttyUSB0") def thread1(): while True: result = serial_device_abstraction.get("foo") print(f"foo查询结果: {result}") time.sleep(1) def thread2(): time.sleep(random.random()) result = serial_device_abstraction.get("bar") print(f"bar查询结果: {result}") threading.Thread(target=thread1).start() threading.Thread(target=thread2).start()
方案2:带优先级调度的进阶版本
如果需要让临时查询(如“bar”)优先执行,避免被循环的“foo”查询长期阻塞,可以用优先级队列+单独工作线程的模式,所有查询请求都进入队列排队,高优先级请求插队执行。
import threading import serial import time import random from queue import PriorityQueue class SerialDeviceAbstraction: def __init__(self, port, baudrate=9600): self.serial_port = serial.Serial(port, baudrate) self._query_queue = PriorityQueue() # 启动单独的工作线程处理串口查询 self._worker_thread = threading.Thread(target=self._process_queries, daemon=True) self._worker_thread.start() def get(self, query, priority=1): # priority值越小,优先级越高(0为最高,1为普通) event = threading.Event() result = [None] # 用列表存储结果,方便在工作线程中修改 self._query_queue.put((priority, query, event, result)) event.wait() # 等待查询完成 return result[0] def _process_queries(self): while True: priority, query, event, result = self._query_queue.get() try: # 实际执行串口读写 self.serial_port.write(query.encode()) data = self.serial_port.read(8) result[0] = data.decode() finally: event.set() # 通知调用线程查询完成 self._query_queue.task_done() # 实例化抽象层 serial_device_abstraction = SerialDeviceAbstraction("/dev/ttyUSB0") def thread1(): while True: result = serial_device_abstraction.get("foo", priority=1) print(f"foo查询结果: {result}") time.sleep(1) def thread2(): time.sleep(random.random()) result = serial_device_abstraction.get("bar", priority=0) # 设置高优先级 print(f"bar查询结果: {result}") threading.Thread(target=thread1).start() threading.Thread(target=thread2).start()
关键说明
- 基础版本:保证线程安全,实现成本低,适合对响应优先级无特殊要求的场景,不会中断正在进行的串口操作,避免数据损坏。
- 进阶版本:通过队列调度实现优先级,高优先级查询可以插队执行,适合需要临时查询优先响应的场景,所有串口操作都在单一线程中执行,彻底规避多线程冲突。
内容的提问来源于stack exchange,提问作者Basj
相关产品推荐
相关产品推荐

