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

串口通信高层抽象实现:多线程并发查询冲突解决问询

串口设备多线程并发查询的通用抽象实现

问题背景

在具备串口通信的硬件设备上,需要持续每秒查询一次“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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 09:48:12