如何避免RTSP断开时Gstream的get_buffer()方法阻塞?
RTSP连接断开时GStreamer
get_buffer()阻塞问题解决方案 问题场景
使用Python+GStreamer开发RTSP监听器POC时,遇到RTSP连接断开后调用get_buffer()会导致程序无限阻塞,无法执行重连逻辑,需要解决阻塞问题并实现自动重连。
原问题代码
import gi gi.require_version('Gst', '1.0') from gi.repository import Gst, GObject, GLib from urllib.parse import quote import cv2 import numpy as np from PIL import Image from io import BytesIO import base64 import time Gst.init(None) passs = quote('asdA@Da32') pipelines = [] rtsp_urls = [ f"rtsp://admin:{passs}@172.16.40.216:7001/21eaf170-115c-f23a-cc1a-0b12e82636dd", f"rtsp://admin:{passs}@172.16.40.216:7001/21eaf170-115c-f24a-cc1a-0b12e82636dd", f"rtsp://admin:{passs}@172.16.40.216:7001/21eaf170-115c-f25a-cc1a-0b12e82636dd", f"rtsp://admin:{passs}@172.16.40.216:7001/21eaf170-115c-f26a-cc1a-0b12e82636dd", ] def retry(): global pipelines for url in rtsp_urls: pipeline = Gst.parse_launch(f'rtspsrc location={url} ! decodebin ! videoconvert ! jpegenc ! appsink name=sink{len(pipelines)}') appsink = pipeline.get_by_name(f'sink{len(pipelines)}') appsink.set_property('sync', False) appsink.set_property('max-buffers', 1) appsink.set_property('drop', True) pipelines.append(pipeline) pipeline.set_state(Gst.State.PLAYING) def rec(): global pipelines print("rec") for pipeline in pipelines: pipeline.set_state(Gst.State.NULL) pipelines.clear() retry() retry() while True: try: for i, url in enumerate(rtsp_urls): start = time.time() try: pipeline = pipelines[i] result, state, pending = pipeline.get_state(0) print(result) print(state) sample = pipeline.get_by_name(f"sink{i}").emit("pull-sample") buffer = sample.get_buffer() except Exception as e: print(e) rec() #retry() if buffer: try: caps = sample.get_caps() height = caps.get_structure(0).get_int('height')[1] width = caps.get_structure(0).get_int('width')[1] data = buffer.extract_dup(0, buffer.get_size()) image = Image.open(BytesIO(data)) print(time.time()-start) image.save(f"{i}zdjecie.jpg", "JPEG") except Exception as e: print(e) else: rec() except Exception as e: print(e)
解决方案
1. 给pull-sample添加超时参数
appsink的pull-sample方法支持传入超时时间(单位:纳秒),默认是无限等待,设置超时后,连接断开时会在超时后抛出异常,触发重连逻辑。
修改点:
# 原代码 sample = pipeline.get_by_name(f"sink{i}").emit("pull-sample") # 修改为(设置500ms超时) sample = pipeline.get_by_name(f"sink{i}").emit("pull-sample", 500 * 1000000)
2. 监听GStreamer总线消息感知连接状态
通过监听GStreamer的总线消息,可实时捕获ERROR或EOS事件,在连接断开时立即触发重连,无需等待pull-sample阻塞。
添加总线监听逻辑:
def bus_call(bus, message, pipeline_idx): t = message.type if t == Gst.MessageType.ERROR: err, debug = message.parse_error() print(f"管道{pipeline_idx}连接错误: {err.message}") rec() elif t == Gst.MessageType.EOS: print(f"管道{pipeline_idx}流结束") rec() return True # 在retry函数中添加总线绑定 def retry(): global pipelines pipelines.clear() for idx, url in enumerate(rtsp_urls): pipeline = Gst.parse_launch(f'rtspsrc location={url} ! decodebin ! videoconvert ! jpegenc ! appsink name=sink{idx}') appsink = pipeline.get_by_name(f"sink{idx}") appsink.set_property('sync', False) appsink.set_property('max-buffers', 1) appsink.set_property('drop', True) # 绑定总线监听 bus = pipeline.get_bus() bus.add_signal_watch() bus.connect("message", bus_call, idx) pipelines.append(pipeline) pipeline.set_state(Gst.State.PLAYING)
3. 替换阻塞轮询为GLib主循环
原代码的while True轮询是阻塞式的,改用GLib主循环处理GStreamer事件,实现非阻塞式运行,同时兼容超时和消息监听逻辑。
修改主循环部分:
# 替换原有的while True loop = GLib.MainLoop() try: loop.run() except KeyboardInterrupt: # 程序终止时清理资源 for pipeline in pipelines: pipeline.set_state(Gst.State.NULL)
综合优化后的完整代码
import gi gi.require_version('Gst', '1.0') from gi.repository import Gst, GObject, GLib from urllib.parse import quote from PIL import Image from io import BytesIO import time Gst.init(None) GObject.threads_init() passs = quote('asdA@Da32') pipelines = [] rtsp_urls = [ f"rtsp://admin:{passs}@172.16.40.216:7001/21eaf170-115c-f23a-cc1a-0b12e82636dd", f"rtsp://admin:{passs}@172.16.40.216:7001/21eaf170-115c-f24a-cc1a-0b12e82636dd", f"rtsp://admin:{passs}@172.16.40.216:7001/21eaf170-115c-f25a-cc1a-0b12e82636dd", f"rtsp://admin:{passs}@172.16.40.216:7001/21eaf170-115c-f26a-cc1a-0b12e82636dd", ] def bus_call(bus, message, pipeline_idx): t = message.type if t == Gst.MessageType.ERROR: err, debug = message.parse_error() print(f"管道{pipeline_idx}连接错误: {err.message}") rec() elif t == Gst.MessageType.EOS: print(f"管道{pipeline_idx}流结束") rec() return True def process_frame(pipeline_idx): start = time.time() try: pipeline = pipelines[pipeline_idx] appsink = pipeline.get_by_name(f"sink{pipeline_idx}") # 设置500ms超时 sample = appsink.emit("pull-sample", 500 * 1000000) if not sample: return True buffer = sample.get_buffer() caps = sample.get_caps() height = caps.get_structure(0).get_int('height')[1] width = caps.get_structure(0).get_int('width')[1] data = buffer.extract_dup(0, buffer.get_size()) image = Image.open(BytesIO(data)) print(f"管道{pipeline_idx}处理耗时: {time.time()-start:.3f}s") image.save(f"{pipeline_idx}zdjecie.jpg", "JPEG") except Exception as e: print(f"管道{pipeline_idx}处理帧错误: {e}") return True def retry(): global pipelines # 清理旧管道资源 for pipeline in pipelines: pipeline.set_state(Gst.State.NULL) pipelines.clear() for idx, url in enumerate(rtsp_urls): pipeline = Gst.parse_launch(f'rtspsrc location={url} ! decodebin ! videoconvert ! jpegenc ! appsink name=sink{idx}') appsink = pipeline.get_by_name(f"sink{idx}") appsink.set_property('sync', False) appsink.set_property('max-buffers', 1) appsink.set_property('drop', True) # 绑定总线监听 bus = pipeline.get_bus() bus.add_signal_watch() bus.connect("message", bus_call, idx) pipelines.append(pipeline) pipeline.set_state(Gst.State.PLAYING) # 定时触发帧处理(每100ms一次) GLib.timeout_add(100, process_frame, idx) def rec(): print("开始重新连接所有RTSP流") retry() if __name__ == "__main__": retry() loop = GLib.MainLoop() try: loop.run() except KeyboardInterrupt: print("程序终止") for pipeline in pipelines: pipeline.set_state(Gst.State.NULL)
内容的提问来源于stack exchange,提问作者Jester48
相关产品推荐
相关产品推荐

