ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

TDengine 流式计算引擎详解:架构、任务编排与状态容错机制

2026/9/14 11:39:01 拓冰建站 浏览量
TDengine 流式计算引擎详解:架构、任务编排与状态容错机制 TDengine 流式计算引擎详解架构、任务编排与状态容错机制【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine本篇基于 TDengine 内部机制文档系统讲解流式计算引擎的完整架构从 SQL 到物理 DAG 的转换流程、mnode 四大逻辑模块的职责、source/agg/sink 三类流任务的分工、检查点与容错设计以及内存管理、流量控制与反压机制。读完本文你可以理解 TDengine 流计算“事件驱动 WAL 数据源 分布式检查点”的设计原理并能结合配置参数如streamBufferSize对集群的流计算性能与资源占用做出针对性调优。流式计算架构总览TDengine 流式计算的完整流程如下当用户输入用于创建流的 SQL 后该 SQL 首先在客户端进行解析生成流式计算执行所需的逻辑执行计划及其相关属性信息随后客户端将这些信息发送至mnodemnode 利用来自数据源超级表所在数据库的 vgroups 信息将逻辑执行计划动态转换为物理执行计划并进一步生成流任务的有向无环图DAG最后mnode启动分布式事务将任务分发至每个 vgroup从而启动整个流式计算流程。从源码结构看流计算的运行时核心实现位于 source/libs/new-stream/ 目录下其中src/stream.c负责任务生命周期管理src/streamReader.c负责从 WAL 读取数据src/streamCheckpoint.c实现检查点逻辑src/streamHb.c实现心跳上报与上文描述的 mnode 各模块一一对应。mnode 包含与流式计算相关的 4 个逻辑模块模块职责任务调度负责将逻辑执行计划转化为物理执行计划并下发到每个 vnodemeta store负责存储流式计算任务的元数据信息以及流任务相应的 DAG 信息检查点调度负责定期生成检查点checkpoint事务并下发到各 source task源任务exec 监控负责接收上报的心跳、更新 mnode 中各任务的执行状态以及定期监控检查点执行状态和 DAG 变动信息此外mnode 还承担着向流式计算任务下发控制命令的重要角色这些命令包括但不限于暂停、恢复执行、删除流任务及更新流任务的上下游信息等。在每个 vnode 上至少部署两个流任务一个是source task负责从 WAL必要时也会从 TSDB中读取数据以供后续任务处理并将处理结果分发给下游任务另一个是sink task写回任务职责是将收到的数据写入所在的 vnode。为了确保数据流的稳定性和系统的可扩展性每个 sink task 都配备了流量控制功能以便根据实际情况调整数据写入速度。核心概念有状态的流式计算流式计算引擎具备强大的标量函数计算能力它处理的数据在时间序列上相互独立无须保留计算的中间状态。对于所有输入数据引擎可以执行固定的变换操作如简单的数值加法并直接输出结果。同时引擎也支持对数据进行聚合计算这类计算需要在执行过程中维护中间状态。以统计设备的日运行时间为例由于统计周期可能跨越多天应用程序必须持续追踪并更新当前的运行状态直至统计周期结束才能得出准确的结果。这正是有状态流式计算的典型场景——引擎在执行过程中必须保持对中间状态的跟踪和管理以确保最终结果的准确性。预写日志WAL当数据写入 TDengine 时首先会被存储在 WAL 文件中。每个 vnode 都拥有自己的 WAL 文件并按照时序数据到达的顺序进行保存。由于 WAL 文件保留了数据到达的顺序它成为流式计算的重要数据来源。此外WAL 文件具有自己的数据保留策略通过数据库参数控制超过保留时长的数据会被从 WAL 文件中清除。这种设计确保了数据的完整性和系统的可靠性同时为流式计算提供了稳定的数据来源。事件驱动事件在系统中指的是状态的变化或转换。在流式计算架构中触发流式计算流程的事件是超级表数据的写入消息。在这一阶段数据可能尚未完全写入 TSDB而是在多个副本之间进行协商并最终达成一致。流式计算采用事件驱动模式执行其数据源并非直接来自 TSDB而是WAL。数据一旦写入 WAL 文件就会被提取出来并加入待处理的队列中等待流式计算任务的进一步处理。这种“数据写入后立即触发计算”的方式确保数据一旦到达就能得到及时处理并能在最短时间内将处理结果存储到目标表中。三种时间概念在流式计算领域时间是一个至关重要的概念。TDengine 的流式计算涉及 3 个关键时间概念事件时间Event Time时序数据中每条记录的主时间戳Primary Timestamp通常由生成数据的传感器或上报数据的网关提供用以精确标记记录的生成时刻。事件时间是流式计算结果更新和推送策略的决定性因素。写入时间Ingestion Time记录被写入数据库的时刻。写入时间与事件时间通常独立一般情况下写入时间晚于或等于事件时间除非出于特定目的用户写入了未来时刻的数据。处理时间Processing Time流式计算引擎开始处理写入 WAL 文件中数据的时间点。对于设置了max_delay选项以控制流式计算结果返回策略的场景处理时间直接决定结果返回的时间。值得注意的是在at_once和window_close这两种计算触发模式下数据一旦到达 WAL 文件就会立即被写入 source task 的输入队列并开始计算。这些时间概念的区分确保流式计算能够准确处理时间序列数据并根据不同时间点的特性采取相应的处理策略。时间窗口聚合TDengine 的流式计算允许根据记录的事件时间将数据划分到不同的时间窗口中通过应用指定的聚合函数计算每个窗口内的聚合结果。当窗口中有新记录到达时系统触发对应窗口的聚合结果更新并根据预先设定的推送策略将更新后的结果传递给下游流式计算任务。当聚合结果需要写入预设的超级表时系统首先根据分组规则生成相应的子表名称然后将结果写入对应的子表。流式计算中的时间窗口划分策略与批量查询中的窗口生成与划分策略保持一致确保了数据处理的一致性和效率。乱序处理在网络传输和数据路由等复杂因素影响下写入数据库的数据可能无法维持事件时间的单调递增特性这种现象被称为乱序写入。乱序写入是否影响相应时间窗口的流式计算结果取决于创建流式计算任务时设置的两个参数watermark水位线控制允许乱序的时间跨度ignore expired忽略过期数据控制是否丢弃超过水位线的过期数据。这两个参数共同作用确定是丢弃这些乱序数据还是将其纳入并增量更新所属时间窗口的计算结果。通过这种方式系统能够在保持流式计算结果准确性的同时灵活处理乱序数据确保数据的完整性和一致性。流式计算任务编排每个激活的流式计算实例都是由分布在不同 vnode 上的多个流任务组成的。这些流任务在整体架构上呈现相似性均包含全内存驻留的输入队列和输出队列、执行时序数据的执行器系统以及用于存储本地状态的存储系统。这种设计确保了流任务的高性能与低延迟同时提供了良好的可扩展性和容错性。按照承担任务的不同流任务可划分为 3 类source task源任务、agg task聚合任务和 sink task写回任务。source task源任务流式计算的数据处理始于本地 WAL 文件中的数据读取这些数据随后在本地节点上进行局部计算。source task 遵循数据到达的自然顺序依次扫描 WAL 文件并筛选出符合特定要求的数据然后对这些时序数据进行顺序处理。因此流式计算的数据源超级表无论分布在多少个 vnode 上集群中都会相应地部署同等数量的源任务。这种分布式处理方式确保了数据的并行处理和高效利用集群资源。agg task聚合任务source task 的下游任务是接收源任务聚合后的结果并对这些结果进行进一步汇总以生成最终输出在集群中配置snode的情况下agg task 会被优先安排在 snode 上执行以利用其存储和处理能力如果集群中没有 snodemnode 则会随机选择一个 vnode在该 vnode 上调度执行 agg task。值得注意agg task 并非在所有情况下都是必需的。对于不涉及窗口聚合的流式计算场景例如仅包含标量运算的流式计算或者数据库只有一个 vnode 时的聚合流式计算就不会出现 agg task此时流式计算的拓扑结构简化为仅包含 source task 和直接输出结果的下游任务两级。sink task写回任务sink task 承担接收 agg task 或 source task 输出结果并将其写入本地 vnode 的重任以此完成数据的写回过程。与 source task 类似每个结果超级表所分布的 vnode 都将配备一个专门的 sink task。用户可以通过配置参数调节 sink task 的吞吐量以满足不同的性能需求。三类任务的协作关系如下source task 的数量直接取决于 vnode 的数量每个 source task 独立负责处理各自 vnode 中的数据与其他 source task 互不干扰不存在顺序性约束而如果最终流式计算结果汇聚到一张表中那么在该表所在的 vnode 上只会部署一个 sink task。流式计算节点 snodesnode是一个专为流式计算服务的独立 taosd 进程专门用于部署 agg task。snode 不仅具备本地状态管理能力还内置了远程备份数据的功能。它使得 snode 能够收集并存储分散在各个 vgroup 中的检查点数据并在需要时将这些数据远程下载到重新启动流式计算的节点上从而确保流式计算状态的一致性和数据的完整性。在仓库中snode 的实现位于 source/dnode/snode/与 bnode、qnode、xnode 等节点类型并列于 dnode 层是集群可选部署的独立组件。状态与容错处理检查点流式计算过程中系统采用分布式检查点机制定期默认每 3 分钟保存计算过程中各任务内部算子的状态这些状态快照即检查点checkpoint。生成检查点的操作仅与处理时间相关联与事件时间无关。在假定所有流任务均正常运行的前提下mnode 定期发起生成检查点的事务并将这些事务分发至每个流的最顶层任务。负责处理这些事务的消息随后会进入数据处理队列。TDengine 的检查点生成方式与业界主流流式计算引擎保持一致每次生成的检查点都包含完整的算子状态信息。副本策略上有两点值得注意对于每个任务检查点数据仅在任务运行的节点上保留一份副本与时序数据存储引擎的副本设置完全独立流任务的元数据信息采用多副本保存机制并被纳入时序数据存储引擎的管理范畴。因此在多副本集群上执行流式计算任务时其元数据信息也将实现多副本冗余。为确保检查点数据的可靠性TDengine 流式计算引擎提供了远程备份检查点数据的功能支持将检查点数据异步上传至远程存储系统。即便本地保留的检查点数据受损也能从远程存储系统下载相应数据并在全新的节点上重启流式计算继续执行计算任务。这一措施进一步增强了流式计算系统的容错能力和数据安全性。从源码结构看检查点的落盘与恢复逻辑集中在 streamCheckpoint.c 中实现心跳上报逻辑在 streamHb.c 中实现二者共同支撑了上文 mnode “检查点调度”与“exec 监控”模块的工作。状态存储后端流任务中算子的计算状态信息以文件的方式持久化存储在本地硬盘中。内存管理流任务的内存占用由以下部分构成项目说明输入/输出队列每个非 sink task 都配备输入队列和输出队列sink task 只有输入队列、不设输出队列。两种队列的数据容量上限均设定为60MB根据实际需求动态占用存储空间队列为空时不占用任何存储空间streamBufferSize控制 agg task 内部保存窗口状态缓存的内存大小默认值 128MBmaxStreamBackendCache限制后端存储在内存中的最大占用存储空间默认同样为128MB文档说明可调范围为16MB 至 1024MB在仓库源码中可以对streamBufferSize的注册和默认值加以印证tglobal.c 中定义了全局变量tsStreamBufferSize单位 MB初始为 0即未显式配置时按总内存的 30% 推算见 tglobal.c并在 tglobal.c 处以128作为缺省值注册为服务端动态配置项与文档所述默认值一致。streamSinkDataRate与maxStreamBackendCache两个参数在文档中有明确记载见下文流量控制部分在当前源码树中未检索到同名配置项使用时以文档说明及发行版实际配置能力为准。流量控制流式计算引擎在 sink task 中实现了流量控制机制以优化数据写入性能并防止资源过度消耗。该机制主要通过以下两个指标控制流量每秒写入操作调用次数sink task 负责将处理结果写入其所属的 vnode该指标的上限被固定为每秒 50 次以确保写入操作的稳定性和系统资源的合理分配每秒写入数据吞吐量通过配置参数streamSinkDataRate控制可调范围为0.1MB/s 至 5MB/s默认值为2MB/s——即对于单个 vnode每个 sink task 每秒最多写入 2MB 数据。sink task 的流量控制机制带来两方面收益一方面能够防止在多副本场景下因高频率写入导致的同步协商缓冲区溢出另一方面避免写入队列中数据堆积过多而消耗大量内存空间有效减少输入队列的内存占用。得益于整体计算框架中应用的反压机制sink task 能够将流量控制的效果直接反馈给最上层的任务从而降低流式计算任务对设备计算资源的占用确保系统整体的稳定性和效率。反压机制TDengine 的流式计算框架部署在多个计算节点上为了协调各节点上的任务执行进度并防止上游任务的数据持续涌入导致下游任务过载系统在任务内部以及上下游任务之间均实施了反压机制。任务内部反压通过监控输出队列的满载情况实现。一旦任务的输出队列达到存储上限当前计算任务便进入等待状态暂停处理新数据直至输出队列有足够空间容纳新的计算结果然后恢复数据处理流程。上下游任务之间反压通过消息传递触发。当下游任务的输入队列达到存储上限即流入下游的数据量持续超过下游任务的最大处理能力时上游任务将接收到下游任务发出的输入队列满载信号此时上游任务适时暂停计算处理直到下游任务处理完毕并允许数据继续分发上游任务才重新开始计算。这种机制有效平衡了任务间的数据流动确保整个流式计算系统的稳定性和高效性。小结TDengine 的流式计算引擎以WAL 为数据源、事件驱动触发通过 mnode 将逻辑计划转换为跨 vgroup 的物理 DAG并借助 source task / agg task / sink task 三级任务拓扑完成分布式计算。其可靠性由分布式检查点默认每 3 分钟、完整算子状态、可选远程备份、snode 的状态汇聚能力以及“60MB 队列上限 streamBufferSize状态缓存 双指标流量控制 队列级反压”的资源管控体系共同保障。如果你希望在集群中启用流计算建议重点关注 snode 的部署、watermark 与 ignore expired 的乱序容忍配置以及streamBufferSize、streamSinkDataRate等参数与业务写入峰值的匹配程度。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考