GStreamer插件开发求助:HTTP Server源元素无法等待数据崩溃问题
问题分析与解决方案
你的核心问题是错误地将libevent服务端逻辑放在了GstPushSrc的create函数中,导致两个致命问题:一是create函数会被GStreamer的推送线程反复调用,重复初始化libevent;二是下游元素启动后立即请求数据,此时无客户端推送数据,create无法返回有效GstBuffer,直接触发流水线崩溃。
下面是具体的修复方案:
1. 分离服务端初始化与数据生产逻辑
不要在create函数中初始化libevent,而是在元素的start虚函数(继承自GstBaseSrc)中完成服务端初始化,确保只执行一次:
- 在元素结构体中维护一个线程安全的队列(推荐用
GAsyncQueue),用来暂存客户端推送的数据块 - 启动单独的线程运行libevent的事件循环,避免阻塞GStreamer的主线程
2. 让create函数阻塞等待数据
GstPushSrc的create函数本身就是同步调用的,你可以利用这一点,让它阻塞直到队列中有数据:
- 调用
g_async_queue_pop()从队列取数据,队列为空时该函数会自动阻塞,正好匹配下游等待数据的需求 - 在元素的
stop方法中向队列推入终止标记,避免流水线停止时create函数无限阻塞
3. 处理HTTP客户端数据并推入队列
在libevent的HTTP请求回调中,将接收到的数据封装为GstBuffer,推入队列:
- 确保设置正确的caps、时间戳等元数据,否则下游元素无法正确处理
- 注意内存管理,推荐用
gst_buffer_new_wrapped()直接包装接收到的数据,避免内存拷贝
核心代码示例
// 自定义元素结构体 typedef struct _HttpSrc { GstPushSrc parent; struct event_base *ev_base; GAsyncQueue *data_queue; gboolean is_stopped; } HttpSrc; // 启动元素时初始化服务端与队列 static gboolean http_src_start(GstBaseSrc *src) { HttpSrc *self = HTTP_SRC(src); self->data_queue = g_async_queue_new(); self->is_stopped = FALSE; // 初始化libevent并启动HTTP服务 self->ev_base = event_base_new(); struct evhttp *http = evhttp_new(self->ev_base); evhttp_bind_socket(http, "0.0.0.0", 8080); evhttp_set_gencb(http, http_request_handler, self); // 启动单独线程运行libevent事件循环 g_thread_new("libevent-loop", (GThreadFunc)event_base_dispatch, self->ev_base); return GST_BASE_SRC_CLASS(parent_class)->start(src); } // 生产数据:阻塞等待队列中的数据 static GstFlowReturn http_src_create(GstPushSrc *src, GstBuffer **buf) { HttpSrc *self = HTTP_SRC(src); gpointer data = g_async_queue_pop(self->data_queue); if (self->is_stopped) { return GST_FLOW_EOS; } *buf = GST_BUFFER(data); return GST_FLOW_OK; } // HTTP客户端请求回调:处理数据并推入队列 static void http_request_handler(struct evhttp_request *req, void *arg) { HttpSrc *self = arg; struct evbuffer *input = evhttp_request_get_input_buffer(req); guint data_len = evbuffer_get_length(input); // 读取客户端推送的数据 guchar *raw_data = g_malloc(data_len); evbuffer_copyout(input, raw_data, data_len); // 封装为GstBuffer(根据实际数据设置caps、时间戳) GstBuffer *buf = gst_buffer_new_wrapped(raw_data, data_len); // gst_buffer_set_caps(buf, gst_caps_from_string("video/x-raw,format=RGB")); // 推入队列,唤醒阻塞的create函数 g_async_queue_push(self->data_queue, buf); // 回复客户端 evhttp_send_reply(req, HTTP_OK, "Received", NULL); } // 停止元素时清理资源 static gboolean http_src_stop(GstBaseSrc *src) { HttpSrc *self = HTTP_SRC(src); self->is_stopped = TRUE; // 推入空指针唤醒阻塞的create函数 g_async_queue_push(self->data_queue, NULL); // 停止libevent事件循环并释放资源 event_base_loopbreak(self->ev_base); event_base_free(self->ev_base); // 清空并销毁队列 g_async_queue_unref(self->data_queue); return GST_BASE_SRC_CLASS(parent_class)->stop(src); }
额外注意事项
- Caps协商:如果客户端推送的音视频格式不固定,需要在接收到第一帧数据后,调用
gst_base_src_set_caps()动态设置caps并通知下游 - 内存安全:确保所有GstBuffer的内存被正确管理,避免泄漏
- 错误处理:在libevent回调中处理客户端连接错误、数据截断等异常情况,避免影响流水线运行
内容的提问来源于stack exchange,提问作者YaleCheung
相关产品推荐
相关产品推荐

