尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

MediaPipe 同步机制深度解析:调度队列、时间戳同步与流控策略的底层原理

MediaPipe 同步机制深度解析:调度队列、时间戳同步与流控策略的底层原理 MediaPipe 同步机制深度解析调度队列、时间戳同步与流控策略的底层原理【免费下载链接】mediapipeCross-platform, customizable ML solutions for live and streaming media.项目地址: https://gitcode.com/GitHub_Trending/med/mediapipe导读本文以官方《Synchronization》概念文档为主线结合本仓库mediapipe的calculator_base.h、input_stream_handler.h、stream_handler/系列实现与flow_limiter_calculator.cc源码系统讲解 MediaPipe 图的三大核心机制——调度队列Scheduler Queue、时间戳同步Timestamp Synchronization与输入策略Input Policy、以及背压与丢包两种流控方案。读完本文你将理解节点何时运行、以怎样的输入集合运行的判定逻辑掌握DefaultInputStreamHandler、SyncSetInputStreamHandler、ImmediateInputStreamHandler的差异与适用场景并能用max_queue_size与FlowLimiterCalculator为实时应用配置可落地的流量控制方案。为什么需要同步无全局时钟的流水线执行MediaPipe 图的执行是去中心化decentralized的整个图没有全局时钟不同节点可以在同一时刻处理不同时间戳的数据。这种设计通过流水线pipelining带来了更高的吞吐量——上游还在解码第 N1 帧时下游已经在渲染第 N 帧的结果。但对于大多数感知perception工作流而言时间信息极其关键。接收多个输入流的节点通常需要协调这些输入。典型例子是目标检测器从某一帧输出一组边界框bounding box这些检测结果必须与原始帧一起送入渲染节点才能画出标注。此时框架就需要保证同一时刻/同一时间戳的数据成组到达。因此MediaPipe 框架的核心职责之一就是为节点提供输入同步。从框架机制上看时间戳timestamp的首要角色正是充当同步键synchronization key。此外MediaPipe 还致力于支持确定性执行在测试、仿真、批处理等场景下至关重要同时允许图作者在需要满足实时约束时放宽确定性。同步与确定性这两个目标直接影响了下文介绍的一系列设计流内时间戳的单调递增约束、timestamp bound 的推进机制、以及各输入策略对就绪ready的判定方式。调度机制一个节点何时被运行数据在 MediaPipe 图中的处理发生在以CalculatorBase子类定义的节点内部。调度系统负责决定每个 calculator 何时运行。调度队列与执行器Executor每个图至少有一个调度队列scheduler queue每个调度队列恰好对应一个执行器executor节点在配置时静态分配到某个队列因此也固定到某个执行器。默认情况下整个图只有一个队列其执行器是一个线程池线程数依据系统能力可用处理器数自动决定。该线程数可通过CalculatorGraphConfig的num_threads字段显式指定见 calculator.proto。如需定制资源使用策略可在图配置的executor列表中声明多个具名执行器字段executor再通过节点上的executor字符串字段把特定节点指派过去——例如让某些低优先级节点跑在低优先级线程上避免它们抢占主链路的计算资源。节点的调度状态与就绪判定每个节点都有一个调度状态可能是not ready未就绪、ready就绪或running运行中。一个**就绪函数readiness function**决定节点是否可以运行它在以下时刻被调用图初始化时某个节点运行结束时某个节点的输入状态发生变化时如有新 packet 入队或 timestamp bound 被推进。在代码层面就绪判定由每个输入流对应的InputStreamHandler完成。基类InputStreamHandler中定义了三态枚举NodeReadinesskNotReady节点尚未就绪kReadyForProcess节点就绪应调用Process()kReadyForClose节点就绪应调用Close()。就绪函数采用哪种判定取决于节点类型源节点source node没有流输入stream input的节点称为源节点它们总是就绪直到它们告知框架不再有数据可输出在Process()中返回tool::StatusStop()见 calculator_base.h 中关于源节点返回值的注释此时源节点被关闭。非源节点当它们有待处理的输入、且这些输入按照该节点**输入策略input policy**所规定的条件构成一个合法输入集合时即为就绪。大多数节点使用默认输入策略少数节点会指定其他策略。注意修改输入策略会改变 calculator 代码对输入所能依赖的保证因此通常不能将使用任意输入策略的 calculator 随意混搭。一个使用特殊输入策略的 calculator 应当专门为其编写并在其 contractGetContract中显式声明。优先级队列就绪任务如何被调度当节点变为就绪一个任务task就会被加入对应的调度队列。调度队列本质上是优先级队列。优先级函数目前是固定的它综合考虑节点自身的静态属性节点在图中的拓扑排序位置。例如更靠近图输出端的节点拥有更高优先级而源节点的优先级最低。这样设计有助于让输出尽早产生、并让流经全图的数据尽快排空。每个队列由各自绑定的 executor 服务executor 负责真正执行任务——即调用 calculator 的Process()代码。正因为执行器可替换、可配置开发者才能在不变更图逻辑的前提下定制计算资源的使用方式。时间戳同步把时间戳当作同步键单调递增约束与 timestamp boundMediaPipe 设计中有一条被同步逻辑依赖的关键不变量推入某一流stream的 packet 的时间戳必须单调递增。这不仅对许多节点是有用的假设更是同步逻辑本身赖以成立的基础。每条流维护一个timestamp bound时间戳下界即允许出现在该流上的新 packet 的最低可能时间戳。当一个时间戳为T的 packet 到达时该流的 bound 会自动推进到T1流内时间戳必须单调递增。这使得框架可以确凿地知道不会再有时戳低于T的 packet 到达。补充实现细节为了让下游节点更早地确定输入状态MediaPipe 还允许流的生产者主动把 bound 推进得比最后一个 packet 暗示的更远提供更紧的 bound例如通过OutputStream::SetNextTimestampBound()。这一能力对降低延迟、配合节流节点工作至关重要FlowLimiterCalculator的源码中就大量使用了它见下文。settled timestamp同步的基石为解释默认输入策略如何工作需要引入settled timestamp已确定时间戳的定义某流上的一个时间戳是settled的当且仅当它小于该流的 timestamp bound。换言之一个时间戳对该流而言已确定意味着该时刻的输入状态已不可更改要么确实存在一个 packet要么可以确定该时间戳的 packet 永远不会到来。进一步推广一个时间戳对多条流而言是 settled 的当且仅当它在每一条流上都是 settled 的如果某个时间戳已 settled那么所有更早的时间戳也必然已 settled。基于这一性质settled 的时间戳可以按升序被确定性地处理。这也解释了为何默认输入策略能够做到绝不乱序、绝不丢包。输入策略Input Policies同步是在每个节点本地完成的使用的是该节点指定的输入策略。框架的调度器/流处理器会在每次有 packet 入队、或某条输入流的下界被推进时触发就绪判定并决定Process()是否被调用。默认策略 DefaultInputStreamHandler默认输入策略由DefaultInputStreamHandler实现当 CalculatorGraph 未显式指定输入流处理器时自动安装见 calculator.proto 图级input_stream_handler字段的注释。它提供确定性的输入同步并保证若多条输入流上出现相同时间戳的 packet无论它们按真实时间到达的先后顺序如何总会被一起处理输入集合按严格升序的时间戳被处理不丢弃任何 packet处理过程完全确定在上述保证的前提下节点尽可能早地进入就绪。推论如果 calculator 在输出 packet 时始终使用当前输入时间戳那么它的输出天然满足单调递增时间戳的要求。警告相反地并不能保证每个输入集合中每条流都有packet 可用——允许某条流在某个时间戳上缺席这正是上面settled定义要表达的确定性问题。结合 settled 的定义默认策略的就绪条件可以精确表述为**存在一个在所有输入流上都已 settled、且在至少一条输入流上包含 packet 的时间戳时calculator 就绪。**该策略把该 settled 时间戳上所有可用的 packet 作为单个input set提供给 calculator。源码 default_input_stream_handler.h 中的注释把这一条件实现为当所有流都结束需要Close()或所有空流上的最小 bound 大于任何流的最小时间戳时意味着下一个时间戳处可能到达的 packet 已经全部收到节点就绪。潜在代价无界的等待与缓冲。确定性行为有一个值得注意的后果对于多输入流的节点等待某个时间戳成为 settled 在理论上可能无限期持续期间缓冲的 packet 数量也可能无界增长。一个典型的病态场景是节点有两个输入流其中一条持续发送 packet而另一条既不发送也不推进 bound于是节点永远等不到那个跨流 settled的时间戳。定制策略一SyncSetInputStreamHandler分组同步为避免上述全有或全无的同步等待框架提供了SyncSetInputStreamHandler它把输入划分成若干独立的同步组sync set每组内部按默认方式各自独立同步。例如一个 calculator 有 5 路输入可以把前 3 路组成一组就像只有这 3 路输入的 calculator 那样同步后 2 路组成另一组。从 头文件注释 可以读出其行为契约calculator每次只会收到来自单个 sync set 的所有可用 packet永远不会同时拿到多个 set 的数据calculator 看到的输入时间戳对每个 sync set 内部是顺序递增的但在不同 set 之间可能跳变节点在任意一个 sync set 就绪时即进入 ready。这非常适合一路高频视频流 一路低频控制/元数据流之类不必严格逐帧对齐的场景。定制策略二ImmediateInputStreamHandler放弃同步另一种极端是ImmediateInputStreamHandler完全不做流间同步任何输入流上一有 packet 就立刻把 calculator 调起来。其头文件特别提醒使用者若不同输入流上连续到达的 packet 具有相同或递减的时间戳该 handler 会以非递增的输入时间戳序列调用 calculator。此时calculator 必须自行负责按所需时间戳累积 packet再进行处理与输出。也就是说确定性责任从框架转移到了 calculator 本身。FlowLimiterCalculator就是它的典型使用者之一——在其GetContract中通过cc-SetInputStreamHandler(ImmediateInputStreamHandler)显式声明见 flow_limiter_calculator.cc。在配置中选用策略除以上三种外本仓库 stream_handler 目录还提供了FixedSizeInputStreamHandler有界缓冲、必要时丢最旧 packet、BarrierInputStreamHandler、EarlyCloseInputStreamHandler、TimestampAlignInputStreamHandler、MuxInputStreamHandler等变体覆盖背压丢帧提前关闭时间戳对齐等更细分的场景。策略可通过两种层级配置见 calculator.proto图级CalculatorGraphConfig.input_stream_handler字段 12对所有未单独指定的节点生效节点级Node.input_stream_handler字段 11覆盖图级设置。框架为节点分配的默认行为正如注释所概括calculator 的Process()在时间戳 t 被调用当且仅当至少一条流在 t 有 packet且其它所有流要么在 t 有 packet、要么已确定不会有即它们的下一个 bound 已大于 t。流量控制Flow Control同步保证了正确性但若生产者的速度远超消费者无限的中间缓冲会带来延迟与内存问题。MediaPipe 提供两套互补的流控机制。机制一背压Backpressure与 max_queue_size第一套机制是背压当某条流上缓冲的 packet 达到可配置的上限时系统会限制上游节点的执行。这个上限由CalculatorGraphConfig::max_queue_size定义。相关字段在 calculator.proto 中的语义如下默认值为 100packet该默认值刻意设置得偏大以保证流水线不被轻易打断将该参数设为-1可完全禁用节流图将按其实际需求使用内存当节点声明了自己的缓冲行为时实际采用max(buffer_size_hint, max_queue_size)机制内建了死锁规避系统当节流导致所有 calculator 都无法运行调度队列为空且无节点在运行时框架会放宽配置的上限防止因错误配置而互相等待。若希望相反——在这种情况下直接让图运行失败以暴露配置错误可设置report_deadlock字段calculator.proto 字段 21。背压机制保持确定性它只做踩刹车不丢任何 packet。代价是当快生产者对慢消费者时图可能整体降速、内存仍可能被上游排队数据撑高。机制二FlowLimiterCalculator 丢包式流控第二套机制是插入特殊节点来丢包丢包依据实时约束进行通常配合自定义输入策略其代表实现是FlowLimiterCalculator。典型拓扑把流控节点放在某子图的输入端同时从子图最终输出连一条回环loopback边回到流控节点在图中声明为back_edge: true的FINISHED输入流。这样流控节点就能持续掌握下游图中正在处理的in-flight时间戳数量并在该数量达到可配置的上限时丢弃新到的 packet。由于 packet 在最上游被丢避免了先处理完一个时间戳、再在中间阶段丢包造成的浪费计算。本仓库真实图中处处可见这一模式例如 face_detection_desktop_live.pbtxt 中的用法node { calculator: FlowLimiterCalculator input_stream: input_video input_stream: FINISHED:output_video input_stream_info: { tag_index: FINISHED back_edge: true } output_stream: throttled_input_video }该图注释清楚说明了它的效果流控节点原样放行第一帧之后等待下游calculator 与子图完成任务才放行下一帧等待期间到达的图像全部丢弃从而把图中大部分区域的 in-flight 图像数限制在 1防止下游节点无限排队。可配置参数定义于 flow_limiter_calculator.proto参数类型/默认值含义max_in_flightint32默认 1同时放行处理的帧数上限max_in_queueint32默认 0等待处理的最大排队帧数超过则丢帧in_flight_timeoutint64 微秒默认 0等待某帧处理完成的超时0 表示不超时超时会放弃过期帧以推进流水线源码注释对调参给出了实践指引见 flow_limiter_calculator.cc 顶部的示例说明max_in_flight: 1max_in_queue: 1通常取得最佳的吞吐/延迟平衡队列里始终有下一帧图不会空闲吞吐接近最优同时队列只保存最新可用帧延迟接近最优将max_in_flight提高到 2 或以上可在图具备高度流水线并行性时换取更高吞吐将max_in_queue降为 0 可改善平均延迟但由于两次输入之间图会短暂空闲吞吐帧率会下降。源码级机制flow_limiter_calculator.cc它在GetContract中声明自身使用ImmediateInputStreamHandler并开启SetProcessTimestampBounds(true)以便感知时间戳下界而非逐个 packetProcess()中通过frames_in_flight_队列与ProcessingAllowed()frames_in_flight_.size() max_in_flight决定放行还是丢弃可选的ALLOW输出流bool packet向外部指示当前时间戳起开始接收/开始丢弃的状态切换。选项既可通过节点node_options静态配置也可在 holistic_tracking_gpu.pbtxt 等图中看到以[type.googleapis.com/mediapipe.FlowLimiterCalculatorOptions]形式内联声明的实例。这套以节点为中心的方案给图作者带来的能力在哪里丢包完全由图控制例如只在输入边界丢、保证下游各中间结果尽量都被消费并允许根据资源约束灵活定制和调整图的实时行为。小结同步与流控的取舍图谱把全文主线收束为一张对照表便于按需选用关注点默认/机制说明就绪判定源节点始终就绪直到返回tool::StatusStop()关闭就绪判定默认输入策略跨所有输入流 settled 且至少一路有 packet 时才处理确定性DefaultInputStreamHandler同戳同处理、严格升序、不丢包、结果确定减等待SyncSetInputStreamHandler分组各自同步组间时间戳可跳变减延迟ImmediateInputStreamHandler有包即处理calculator 自担累积/顺序职责背压max_queue_size默认 100-1 关闭限速不丢包内建死锁规避可report_deadlock报错实时丢帧FlowLimiterCalculator配合FINISHED回环边限制 in-flight 帧数从源头丢包同步机制保证了确定性这一 MediaPipe 的重要特性而时间戳单调约束 timestamp bound settled 判定是理解整个框架的钥匙。若希望进一步深入可在本仓库中继续阅读相关文档与实现框架概念graphs.md、packets.md、calculators.md、realtime_streams.md底层实现scheduler_queue.h、input_stream_handler.h、executor.h各策略源码stream_handler 目录每个 handler 均附有对应测试如 default_input_stream_handler_test.cc。【免费下载链接】mediapipeCross-platform, customizable ML solutions for live and streaming media.项目地址: https://gitcode.com/GitHub_Trending/med/mediapipe创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表