公司动态
GstQueue 知识文档
GstQueue 知识文档GStreamer 版本1.0subprojects/gstreamer1. 概述queue是 GStreamer 核心插件中最基础的缓冲队列元素ELEMENT需要状态催动。它在管道中承担两个核心职责职责说明线程解耦queue 在 sink pad入与 src pad出之间建立一条独立线程边界。上游在调用线程里chain入队下游由 queue 自己的loop线程出队推送。缓冲与流控通过缓冲一定量的数据吸收上下游速率差异并在队列满/空时施加反压或丢弃策略。一句话queue 是一个带阈值控制的线程安全 FIFO把管道切成两个独立运行的线程域。2. 调试与设计时最该关注的技术点2.1 先确认 queue 是否真正建立了线程边界queue不是一个普通的缓存结构而是一个GstElement。只有 pad 激活并进入GST_PAD_MODE_PUSH后src pad 才会启动gst_queue_looptask数据才会由独立的 src 线程出队并推向下游。重点代码类型与注册G_DEFINE_TYPE_WITH_CODE(..., GST_TYPE_ELEMENT, ...)、GST_ELEMENT_REGISTER_DEFINE(..., queue, ...)src 线程启动gst_queue_src_activate_mode()src 线程入口gst_queue_loop()调试时不要只看 queue 的状态要确认srcpad是否 activegst_pad_start_task()是否成功上游是否持续进入gst_queue_chain()下游是否阻塞在gst_pad_push()。如果 src task 没启动队列只会不断入队最终表现为“queue 满后上游卡住”。2.2 满队列时默认是反压不是丢数据gst_queue_is_filled()对max-size-buffers、max-size-bytes、max-size-time使用或关系任一非零上限达到就判定为 full。默认leaky GST_QUEUE_NO_LEAK入队线程会在 gst_queue_chain_buffer_or_list() 中执行g_cond_wait(item_del, queue-qlock)直到 src 线程出队并发出item_del信号。因此排查“下载停止”时要区分下游不消费 - queue src 线程在 gst_pad_push() 阻塞 - queue 不再出队 - queue sink 线程在 item_del 上等待 - 上游 gst_pad_push() 阻塞 - source 不再继续读取/下载确认是否丢数据首先检查leakyno阻塞不丢upstream丢新 bufferdownstream丢队头旧 buffer。2.3 关注qlock和条件变量不要把“推下游”误认为“持锁读队列”GstQueue用queue-qlock保护GstQueueArray、水位统计和流控状态item_add入队后唤醒等待数据的 src loopitem_del出队后唤醒因满而阻塞的 sink chainquery_handled协调序列化 query。src 线程从队列取出 item 后会先解锁再调用gst_pad_push()推送返回后重新加锁并检查srcresult。这不是遗漏锁而是避免下游阻塞时占用qlock否则上游无法入队甚至可能形成锁依赖问题。排查卡死时重点看三个等待点GST_QUEUE_WAIT_ADD_CHECK队列空src 线程等待数据GST_QUEUE_WAIT_DEL_CHECK队列满sink 线程等待空间gst_pad_push()queue 已出队但下游处理或反压阻塞。2.4 阈值是固定配置queue 不会自适应扩容max-size-*和min-threshold-*是普通 GObject 属性。queue 只会读取当前值与cur_level比较不会依据码率、下载速度或播放状态自动调整。设计时必须明确max-size-*决定什么时候满min-threshold-*决定什么时候允许继续出队current-level-*只反映当前水位不是控制参数若要动态调整必须由外部代码调用g_object_set()或使用queue2/multiqueue的专门能力。2.5 状态切换与 flush 必须成对观察queue 的线程不是永久运行的。pad deactivate、FLUSH_START、下游 flow error 都可能设置srcresult为非GST_FLOW_OK唤醒等待线程并暂停 task。关键流程FLUSH_START - 转发 flush event - srcresult GST_FLOW_FLUSHING - signal item_add / item_del - 停止 src task - 等待线程退出 FLUSH_STOP - 清空队列并恢复水位统计 - srcresult GST_FLOW_OK - 清除 EOS/异常标志 - 重新 start src task如果只调用停止或只发送 flush 的一半常见结果是chain 线程仍睡在item_delloop 线程仍睡在item_addseek 后旧 buffer 没清掉queue 永久返回GST_FLOW_FLUSHING。2.6 数据完整性不只取决于 buffer还取决于事件顺序序列化事件例如STREAM_START、SEGMENT、TAG、EOS会和 buffer 一起进入 FIFO保证下游看到的顺序非序列化事件则直接转发。flush 时 queue 会清理缓存 item但会保留必要的 sticky eventSEGMENT和EOS有特殊处理。排查 seek、切流、EOS 后首帧异常时应同时检查src_segment/sink_segmentsrcresult、eos、unexpected是否有旧的SEGMENT或 sticky event 被清理/重新发送首个新 buffer 是否需要GST_BUFFER_FLAG_DISCONT。2.7 实际调试优先观察哪些量建议同时观察 queue 属性和 GStreamer debug 日志gst-launch-1.0-m\...!queuenameq\21|GST_DEBUGqueue:6,queue_dataflow:6...重点属性current-level-buffers current-level-bytes current-level-time max-size-buffers max-size-bytes max-size-time min-threshold-buffers min-threshold-bytes min-threshold-time leaky重点信号underrun、running、overrun、pushing。注意这些信号在 streaming thread 中触发回调里不要执行耗时操作或再次同步等待同一条管道。3. 线程模型上游线程 queue 元素 下游线程queue 自建 │ ┌──────────┐ ├─ gst_pad_push ─────────▶ │ sinkpad │ │ (chain) │ chain │─┐ │ └──────────┘ │ enqueue加锁 │ ▼ │ ┌─────────────────┐ │ │ GstQueueArray │ FIFO 环形数组 │ │ (queue-queue) │ │ └─────────────────┘ │ │ dequeue加锁 │ ┌──────────┐ │ │ │ srcpad │◀┘ │ │ loop │─────────▶ gst_pad_push 到下游 │ └──────────┘入队侧上游调用gst_queue_chain→gst_queue_chain_buffer_or_list运行在上游的流线程。出队侧gst_queue_loop运行在srcpad 自己的 task 线程由gst_queue_push_one取数据推给下游。两侧通过queue-qlockGMutex互斥通过两个条件变量协调item_add队列由空变非空时唤醒等待出队的 loop 线程。item_del队列由满变不满时唤醒被反压阻塞的 chain 线程。条件变量的等待/唤醒被封装成宏GST_QUEUE_WAIT_ADD_CHECK/GST_QUEUE_WAIT_DEL_CHECK/GST_QUEUE_SIGNAL_ADD/GST_QUEUE_SIGNAL_DEL。每次从等待中醒来都会检查srcresult非GST_FLOW_OK如 flush就跳到错误处理实现干净退出。3.1GstQueue的核心私有状态GstQueue定义在gstqueue.h关键字段如下省略类型注册和 TCL 条件编译字段struct_GstQueue{GstElement element;GstPad*sinkpad;GstPad*srcpad;GstSegment sink_segment;GstSegment src_segment;GstClockTimeDiff sinktime,srctime;gboolean sink_tainted,src_tainted;GstFlowReturn srcresult;gboolean unexpected;gboolean eos;GstQueueArray*queue;GstQueueSize cur_level;GstQueueSize max_size;GstQueueSize min_threshold;GstQueueSize orig_min_threshold;gint leaky;GMutex qlock;gboolean waiting_add;GCond item_add;gboolean waiting_del;GCond item_del;GCond query_handled;};字段分组理解字段作用调试价值queue真正保存GstBuffer/GstBufferList/ 序列化 event / query 的 FIFO判断是否真的有缓存 item而非只看播放器状态cur_level当前 buffers/bytes/time 水位判断卡在“空等”还是“满等”max_size、min_threshold固定的满/可读阈值判断反压与出队迟滞的原因srcresultsrc task 的最近 flow 返回值GST_FLOW_FLUSHING、GST_FLOW_EOS、GST_FLOW_NOT_LINKED会改变所有等待线程的退出路径eos、unexpected拒收后续数据或标记异常流控排查 EOS 后无法继续 pushsink_segment、src_segment、sinktime、srctime按 segment 计算current-level-time排查 timestamp/seek 后 time 水位异常qlock、item_add、item_del保护共享状态并协调入/出队排查死锁、反压和 task 唤醒3.2 入队、出队与水位更新的实际调用链入队上游流线程gst_pad_push(upstream src pad, buffer) - gst_queue_chain() - gst_queue_chain_buffer_or_list() - 持有 qlock检查 srcresult/eos/unexpected - 若 full按 leaky 策略丢弃或等待 item_del - gst_queue_locked_enqueue_buffer() - cur_level.buffers - cur_level.bytes gst_buffer_get_size(buffer) - apply_buffer(... sink_segment ...) 更新 cur_level.time - gst_queue_array_push_tail_struct(queue-queue, ...) - GST_QUEUE_SIGNAL_ADD() 唤醒 src loop - GST_FLOW_OKgst_queue_locked_enqueue_buffer()只在已经持有qlock时调用它先更新水位再把带GstMiniObject *item与字节大小的GstQueueItem追加到 FIFO 尾部。出队queue 的 src task 线程gst_queue_loop() - 持有 qlock - 若 gst_queue_is_empty()等待 item_add - gst_queue_push_one() - gst_queue_locked_dequeue() 从 FIFO 头取一个 item - 先释放 qlock - gst_pad_push(queue-srcpad, buffer) // 同步调用下游 chain - 重新持有 qlock检查 srcresult / push 返回值 - 保存 queue-srcresult重要细节gst_pad_push()前 buffer 已从GstQueueArray移除queue 对该 buffer 持有的引用已经转交给局部变量/下游故而 push 期间允许释放qlock。这让上游可继续入队直到 queue 真正达到 max-size而不是被一次慢下游 push 无条件串行化。3.3 判满/判空的真实代码判满是三个最大阈值的或关系staticgbooleangst_queue_is_filled(GstQueue*queue){return(queue-max_size.buffers0queue-cur_level.buffersqueue-max_size.buffers)||(queue-max_size.bytes0queue-cur_level.bytesqueue-max_size.bytes)||(queue-max_size.time0queue-cur_level.timequeue-max_size.time);}判空不是简单的GstQueueArray为空FIFO 没有 item 时当然为空FIFO 尾部是 event/query 时视为非空避免序列化 query 永远等不到处理尾部是 buffer 时如果任一启用的min-threshold-*尚未达到且队列又没有满也视为“暂不可读”“未达到 min 但已满”必须改为可读否则上游在满等待、下游在空等待形成永久死锁。这个“min threshold未达但max size已满时必须放行出队”的保护是设计 max/min 阈值时最容易忽略的边界条件。3.4 Pad 激活、停止和 task 生命周期src pad 的 push-mode activation 是 queue 工作的开关if(active){queue-srcresultGST_FLOW_OK;queue-eosFALSE;queue-unexpectedFALSE;gst_pad_start_task(pad,gst_queue_loop,pad,NULL);}else{queue-srcresultGST_FLOW_FLUSHING;g_cond_signal(queue-item_add);// 叫醒空队列等待的 loopgst_pad_stop_task(pad);gst_queue_locked_flush(queue,FALSE);}sink pad deactivate 则采用相反顺序先置GST_FLOW_FLUSHING并唤醒item_del让可能正在满队列等待的 chain 返回再拿GST_PAD_STREAM_LOCK等上游流线程退出最后清空 FIFO。这个顺序避免 flush/deactivate 与仍在执行的 chain 并发释放同一 item。4. 三个维度的容量控制queue 同时用三个维度衡量满任一维度达到上限即视为满0 表示该维度不限制queue 用三个维度衡量数据量每个维度都有一组最大 / 最小 / 当前属性。三者互相独立同时生效。所有阈值都是固定的 GObject 属性——一旦设定queue 内部永不主动改动详见 §4.4。4.1 最大阈值 max-size-*判满任一维度达到上限即视为满触发 leaky/反压逻辑。设为0表示该维度不参与判断。维度属性类型取值范围默认值宏含义缓冲个数max-size-buffersuint0 ~G_MAXUINT200DEFAULT_MAX_SIZE_BUFFERS队列中 buffer 数量上限字节数max-size-bytesuint0 ~G_MAXUINT10 MBDEFAULT_MAX_SIZE_BYTES队列中数据总字节上限时间max-size-timeuint640 ~G_MAXUINT641 秒GST_SECOND 1e9 nsDEFAULT_MAX_SIZE_TIME队列中数据的时间跨度上限纳秒属性标志READWRITE | MUTABLE_PLAYING——可在 PLAYING 状态下动态修改由外部调不是 queue 自己调。初值在gst_queue_init里从上述DEFAULT_*宏赋给queue-max_size.{buffers,bytes,time}。4.2 最小阈值 min-threshold-*判可读/滞后用于滞后hysteresis队列数据量低于最小阈值时即使非空loop 线程也不出队直到重新积累到阈值以上才恢复出队。避免下游被一点一点挤牙膏常用于需要成块数据的场景。默认全0关闭即非空就能读。维度属性类型取值范围默认值含义缓冲个数min-threshold-buffersuint0 ~G_MAXUINT0禁用允许出队所需的最小 buffer 数字节数min-threshold-bytesuint0 ~G_MAXUINT0禁用允许出队所需的最小字节数时间min-threshold-timeuint640 ~G_MAXUINT640禁用允许出队所需的最小时间跨度纳秒初值在gst_queue_init里由GST_QUEUE_CLEAR_LEVEL清零queue-min_threshold与queue-orig_min_threshold均置 0。同样是READWRITE | MUTABLE_PLAYING。4.3 当前水位 current-level-*只读实时反映队列当前占用供监控/调试用。维度属性类型权限含义缓冲个数current-level-buffersuint只读当前 buffer 数量字节数current-level-bytesuint只读当前数据字节数时间current-level-timeuint64只读当前数据时间跨度纳秒对应内部字段queue-cur_level.{buffers,bytes,time}初始由GST_QUEUE_CLEAR_LEVEL清零。4.4 阈值是固定的——queue 不做任何自适应这三组阈值中的 max/min 一旦设定queue 内部从不主动调整。gst_queue_is_filled()/gst_queue_is_empty()每次都是拿实时cur_level去比这几个固定的 max/min 值没有扩容、没有根据码率/带宽/交织度自适应的逻辑。queue 的自我调节能力为零只有满了就阻塞/丢、低于 min 就不读这几个固定反应。§6 提到的 overrun 信号回调可能改阈值指的是外部代码监听信号后主动g_object_set改不是 queue 自己改。若要真正的动态自适应阈值应使用multiqueue按流交织度动态调整时间上限、饥饿时扩展深度。4.5 判定逻辑判满gst_queue_is_filled()——任一 max 维度非 0被cur_level达到即为满。判空gst_queue_is_empty()——不仅要队列无数据还要考虑 min-threshold低于阈值也视为不可读等效于对出队而言的空。5. Leaky 策略满了怎么办——重点这是 queue 最容易被误解的部分。当队列已满、上游又要 push 新数据时行为由leaky属性决定。枚举定义queue_leaky_get_type枚举值nick行为GST_QUEUE_NO_LEAKno默认。不丢数据阻塞上游线程直到腾出空间GST_QUEUE_LEAK_UPSTREAMupstream丢弃新来的那个 buffer队尾并给下一个 buffer 打DISCONTGST_QUEUE_LEAK_DOWNSTREAMdownstream丢弃队头最老的 buffer给新 buffer 腾位默认值GST_QUEUE_NO_LEAK在属性安装处确认g_object_class_install_property(gobject_class,PROP_LEAKY,g_param_spec_enum(leaky,Leaky,Where the queue leaks, if at all,GST_TYPE_QUEUE_LEAKY,GST_QUEUE_NO_LEAK,/* ← 默认不丢 */...));5.1 核心逻辑gst_queue_chain_buffer_or_listwhile(gst_queue_is_filled(queue)){/* 满了先发 overrun 信号非 silent 时信号处理里可能改阈值 */if(!queue-silent){g_signal_emit(queue,gst_queue_signals[SIGNAL_OVERRUN],0);if(!gst_queue_is_filled(queue))/* 阈值被改大后重新判断 */break;}switch(queue-leaky){caseGST_QUEUE_LEAK_UPSTREAM:queue-tail_needs_discontTRUE;gotoout_unref;/* 直接丢掉新 buffer */caseGST_QUEUE_LEAK_DOWNSTREAM:gst_queue_leak_downstream(queue);/* 丢队头循环再判 */break;caseGST_QUEUE_NO_LEAK:default:/* 不丢阻塞等待出队信号直到不再满 */while(gst_queue_is_filled(queue)){GST_QUEUE_WAIT_DEL_CHECK(queue,out_flushing);}break;}}要点NO_LEAK 是阻塞而非丢弃。GST_QUEUE_WAIT_DEL_CHECK内部g_cond_wait(item_del, qlock)让出上游线程直到 loop 出队后 signal。这就是 GStreamer反压backpressure的实现根基——队列满 → 阻塞上游 chain → 上游的gst_pad_push返回不了 → 逐级向上游传导最终掐停 source如souphttpsrc停止读 socketTCP 反压到服务器。overrun 信号先于丢弃/阻塞发出且发信号时会先解锁再回锁信号回调有机会临时调大阈值gst_queue_is_filled重判从而避免真正丢弃/阻塞。DISCONT 标记无论上游漏还是下游漏都会给接续的那个 buffer打上GST_BUFFER_FLAG_DISCONT通知下游数据不连续tail_needs_discont/head_needs_discont。5.2 何时选哪种场景建议普通点播/文件播放不能丢数据no默认靠反压保证完整性实时源摄像头、直播采集宁可丢帧也不能卡downstream丢老数据保时新只关心最新态、旧帧无意义downstream需要限制上游过快但又不想阻塞少见upstream6. 信号queue提供四个信号silentTRUE时全部不发以降低开销信号触发时机典型用途underrun队列变空监控下游消费过快running从满恢复到不满NO_LEAK 阻塞解除后与 overrun 配对overrun队列变满动态调阈值 / 统计pushing开始向下游推送调试信号发射一律遵循先解锁 → emit → 回锁并重查srcresult的模式避免在持锁状态下调用用户回调导致死锁。7. 事件与查询的处理序列化事件如SEGMENT、EOS、TAG随数据一起入队gst_queue_locked_enqueue_event保持与 buffer 的先后顺序。非序列化事件如FLUSH_START、QOS不入队立即gst_pad_push_event直通下游。EOS入队后置queue-eos此后 chain 拒收任何新数据goto out_eos返回GST_FLOW_EOS。flush-on-eos默认FALSE置TRUE时收到 EOS 直接清空队列并尽快转发 EOS适合直播采集想快速收尾的场景可能丢队尾数据。序列化查询也会入队chain 线程g_cond_wait(query_handled)等 loop 线程处理完再返回结果last_handled_query保证查询看到的是队列排空后的下游状态。8. Flush 流程flush 时FLUSH_START置srcresult GST_FLOW_FLUSHING。唤醒所有条件变量等待者item_add/item_del/query_handled令阻塞在WAIT_ADD/WAIT_DEL的线程立即通过srcresult检查跳到out_flushing干净退出。gst_queue_locked_flush清空GstQueueArray中残留 item 并 unref。FLUSH_STOP后重置srcresult GST_FLOW_OK恢复正常收发。这套机制保证 flush如 seek时不会有线程卡在满/空等待上。9. 与 queue2 / multiqueue 的区别元素定位queue轻量线程边界 简单 FIFO leaky 策略。纯内存无磁盘缓存无 buffering 消息只有信号。queue2面向网络/大缓冲支持磁盘环形缓冲download buffering、发BUFFERING消息、带 buffering 百分比。multiqueue多路流各一条 SingleQueue解决多流饥饿与 non-linked 竞速playbin/decodebin 内部使用。选择原则单纯要个线程解耦或小缓冲用queue要网络缓冲/进度条用queue2解复用后多路流用multiqueue。10. 常见误区澄清“队列满了会丢数据”—— 默认leakyno不丢而是阻塞上游形成反压。只有显式设leakyupstream/downstream才丢。“反压会无限下载”—— 恰恰相反。NO_LEAK 阻塞上游反压逐级传导到 source下载会在缓冲填满后停住不会无限下载。“queue 会做 buffering 进度提示”—— 不会。那是queue2的能力queue只有信号没有BUFFERING消息。“三个 max-size 是与关系还是或关系”—— 或关系任一维度达上限即满。设为 0 的维度不参与判断。11. 关键函数索引函数作用gst_queue_chain/gst_queue_chain_buffer_or_list入队主逻辑含 leaky/反压判断gst_queue_loop/gst_queue_push_one出队并推送下游gst_queue_is_filled/gst_queue_is_empty满/空判定含 min-thresholdgst_queue_locked_enqueue_buffer/_event/_dequeue加锁的入/出队原语gst_queue_leak_downstream下游漏丢队头gst_queue_handle_sink_event事件序列化/直通分发gst_queue_locked_flushflush 时清空队列