Python实现Google Coral目标检测数据向多消费者非阻塞异步传输的方法
问题原因
你原来的写法有两个核心问题导致达不到预期:
- Python普通生成器是一次性可迭代对象,只要被第一个消费者迭代消费完,后面的消费者再迭代就拿不到任何数据了
- 同步迭代生成器的逻辑是阻塞的,调用
VideoViewer().set_stream(generator)的时候会直接进入for循环迭代生成器,永远执行不到后面给HttpApi设置流的代码
解决方案
对新手来说最容易上手的方案是用「线程+队列」实现广播分发,不需要改你原来的检测生成器代码,也不需要复杂的异步逻辑,核心思路是:
- 用一个独立线程跑检测生成器,每产出一条结果就推送到所有消费者的专属队列里
- 每个消费者各自开独立线程,从自己的队列里拿数据处理,互不影响,也不会阻塞主线程
完整实现代码如下:
import cv2 import threading import queue from PIL import Image from cv2 import VideoCapture # 你原来的Detector类保留即可,这里只做示例 class Detector: def _detect_image(self, img): # 你原来的检测逻辑 return [{"bbox": [10,10,50,50], "class": "person"}] def detect_with_video(self, video: VideoCapture): """Detect objects in a video""" streaming = True frame_id = 0 results = [] detection_delay = 4 while streaming: _, frame = video.read() if frame is None: break if frame_id % detection_delay == 0: results = self._detect_image(Image.fromarray(frame)) for result in results: yield result frame_id +=1 # 流广播器:负责把生成器的结果分发给所有订阅的消费者 class StreamBroadcaster: def __init__(self, generator, max_queue_size=10): self.generator = generator self.max_queue_size = max_queue_size # 队列最大长度,避免处理太慢爆内存 self.consumers = [] # 启动后台线程跑检测生成器 self.broadcast_thread = threading.Thread(target=self._broadcast_loop, daemon=True) self.broadcast_thread.start() def _broadcast_loop(self): for item in self.generator: # 给每个消费者的队列推送同一份数据 for q in self.consumers: try: q.put_nowait(item) except queue.Full: # 队列满了就丢弃最老的一条,保证实时性,不需要丢数据可以删掉这个逻辑改成q.put(item) q.get_nowait() q.put_nowait(item) def subscribe(self): # 新消费者订阅,返回专属队列 new_queue = queue.Queue(maxsize=self.max_queue_size) self.consumers.append(new_queue) return new_queue class VideoViewer: """VideoViewer launches video windows""" def new_window(self): # 你原来的新建窗口逻辑 pass def set_stream(self, queue): # 开独立线程处理流,不阻塞主线程 self.work_thread = threading.Thread(target=self._render_loop, args=(queue,), daemon=True) self.work_thread.start() def _render_loop(self, q): while True: item = q.get() # 替换成你原来的画边界框、渲染窗口逻辑 print(f"VideoViewer收到结果:{item}") class HttpApi: def set_stream(self, queue): # 开独立线程处理流,不阻塞主线程 self.work_thread = threading.Thread(target=self._push_loop, args=(queue,), daemon=True) self.work_thread.start() def _push_loop(self, q): while True: item = q.get() # 后续websocket推送逻辑写在这里 print(f"HttpApi收到结果:{item}") class App: """The application""" def __init__(self): self.video: VideoCapture = cv2.VideoCapture(0) self.detector = Detector() def run(self): # 初始化视频窗口 viewer = VideoViewer() viewer.new_window() # 创建广播器 detect_generator = self.detector.detect_with_video(self.video) broadcaster = StreamBroadcaster(detect_generator) # 两个消费者分别订阅,拿到专属队列 viewer_queue = broadcaster.subscribe() api_queue = broadcaster.subscribe() # 分别设置流,都是非阻塞调用 viewer.set_stream(viewer_queue) HttpApi().set_stream(api_queue) # 主线程阻塞等待退出信号,避免程序直接结束 input("按任意键退出\n") self.video.release() if __name__ == "__main__": App().run()
注意事项
- 代码里用的
daemon=True标记的线程会随着主程序退出自动销毁,不需要手动管理生命周期 - 如果你的场景不能丢检测数据,删掉队列满了丢旧数据的逻辑,改成普通阻塞
put即可 - 后续如果要对接asyncio实现的websocket接口,把普通
queue.Queue换成asyncio.Queue,广播逻辑换成协程实现即可,整体思路完全一致
内容的提问来源于stack exchange,提问作者Tom
相关产品推荐
相关产品推荐

