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

如何避免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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 08:42:42