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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 19:54:28