TDengine 流计算可观测性与故障排查实战:基于 ins_streams / ins_stream_tasks / ins_stream_recalculates 系统视图
TDengine 流计算可观测性与故障排查实战基于 ins_streams / ins_stream_tasks / ins_stream_recalculates 系统视图【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengineTDengine 将流计算Stream Processing任务的运行状态以系统视图的形式暴露在information_schema数据库中。本文围绕04-observability.md介绍的三张核心视图——ins_streams流级总览、ins_stream_tasks任务级下钻、ins_stream_recalculates手动重算任务——系统讲解如何监控实时处理延迟、吞吐、历史数据处理进度与重算进度并提供一套可直接落地的故障排查流程。读完本文你将能独立定位哪个流、哪个节点、哪个任务出现了问题并正确解读这些指标中容易误判的NULL与零值语义。一、三张系统视图的分工与定位思路TDengine 的流计算在运行时会被拆分为多种内部任务Task分布在集群的snode上执行。information_schema中的三张视图正是围绕这一架构设计的视图定位层级主要用途information_schema.ins_streams流Stream级快速确认流的整体状态、错误信息、实时延迟、吞吐与历史进度information_schema.ins_stream_tasks任务Task级定位具体节点上的任务查看任务健康度与指标新鲜度information_schema.ins_stream_recalculates重算Recalculation级查看手动重算任务的时间范围、进度与状态推荐的定位思路是从粗到细先用流级视图发现问题再下钻到任务级视图定位节点或任务最后在需要时检查重算任务。这三张视图的完整列定义可以在 01-meta.md 中查到它们也分别与SHOW STREAMS、SHOW STREAM TASKS等SHOW命令等价但相比SHOW语句SELECT ... FROM information_schema.xxx支持过滤、排序以及任意 TDengineSELECT能力更适合做监控与告警。从源码结构看这些视图由管理节点mnode侧的 mndStream.c、mndStreamMgmt.c、mndShow.c 等模块维护并通过 systable.c 注册到information_schema中任务运行期的一分钟级指标则由 streamTaskStats.c 采集与汇总。二、查看流整体状态ins_streams以下查询展示了流的整体状态、错误信息以及核心运行指标SELECT stream_name, status, message, realtime_lag_ms, input_rows_per_sec_1m, output_rows_per_sec_1m, runner_result_latency_avg_1m_ms, history_progress_pct FROM information_schema.ins_streams ORDER BY stream_name;各核心指标的含义如下Metric单位说明realtime_lag_ms毫秒所有有效入口 Reader 中最慢者的实时处理延迟input_rows_per_sec_1m行/秒最近一个完整 60 秒窗口内Trigger 接受的逻辑输入速率output_rows_per_sec_1m行/秒最近一个完整 60 秒窗口内所有最终结果 Runner 成功交付的结果速率runner_result_latency_avg_1m_ms毫秒最近一个完整 60 秒窗口内Runner 从开始处理一次计算请求到形成逻辑结果的加权平均耗时history_progress_pct百分比建流时配置的历史数据处理完成进度取值范围 0100解读这几个指标时有三点需要特别注意realtime_lag_ms取所有有效入口 Reader 中最慢者。已经追上进度、但暂时没有新数据的 Reader 不会让该值无限增长。如果流引用了外部数据源则没有 WAL 进度该字段恒为NULL。输入速率与输出速率属于不同数据层不能直接相减计算丢失率。输入速率统计的是经过过滤和路由后、被流接受的逻辑行输出速率只统计成功交付的最终结果行。窗口聚合、过滤、以及一行输入产生多行结果都会导致两者数值不一致。结果延迟的统计口径。runner_result_latency_avg_1m_ms只统计 Runner 从开始处理计算请求到形成逻辑结果的时间不含请求到达 Runner 之前的排队与网络时间也不含结果形成之后的结果表写入或通知发送时间。三、下钻到任务级ins_stream_tasks当流状态异常、流级指标为NULL或者需要定位具体节点时查询任务视图SELECT stream_name, task_id, type, deploy_id, node_type, node_id, status, last_update, message, input_rows_per_sec_1m, output_rows_per_sec_1m, runner_result_latency_avg_1m_ms FROM information_schema.ins_stream_tasks WHERE stream_name your_stream_name ORDER BY type, deploy_id, task_id;任务级指标的可用性取决于任务职责INS_STREAM_TASKS入口 Reader提供物理输入速率统计该 Reader 实际读取并处理的行数负责最终结果交付的 Runner提供输出速率和结果形成延迟Trigger 任务、计算数据 Reader、非最终结果 Runner上述指标列均为NULL。同时用status判断任务是否健康用last_update判断其状态与指标是否仍然新鲜。注意流级视图没有last_update列新鲜度判断只能依赖任务级视图。从实现上看任务指标与流级指标由 streamTaskStats.c 统一采集并上抛给 mnode 汇总因此任务级与流级指标在口径上保持一致——流级输入速率对应入口 Reader 的输入流级输出速率对应最终结果 Runner 的输出。四、查看历史数据处理进度history_progress_pct对于使用STREAM_OPTIONS(FILL_HISTORY)或STREAM_OPTIONS(FILL_HISTORY_FIRST)创建的流具体语法见 01-syntax.mdins_streams.history_progress_pct会报告初始历史时间范围内的处理进度099历史数据处理尚未完成100历史数据处理已完成NULL未启用历史数据处理或当前没有有效的进度值。需要特别注意的是该百分比表示原始历史时间范围的覆盖率而不是已经产出的结果行数占比。FILL_HISTORY与FILL_HISTORY_FIRST的区别在于前者允许历史处理与实时处理并行推进后者要求严格按时间顺序优先处理历史数据、历史处理完成前不开始实时计算。两者不可同时使用且均不支持PERIOD定时触发模式。五、查看手动重算任务ins_stream_recalculates当需要针对某个时间范围重新计算时可通过RECALCULATE STREAM发起手动重算该语法在 sql.y 中有对应文法支持。以下查询展示每个手动重算任务的时间范围、进度与状态SELECT stream_name, recalc_id, start, end, progress, status, request_time, message FROM information_schema.ins_stream_recalculates WHERE stream_name your_stream_name ORDER BY start, recalc_id;重算状态机Status说明Pending请求已被接受但重算尚未开始Running重算已开始但尚未完成Finished重算完成progress为100%Failed发生了不可恢复的错误重算无法完成关键语义滚动升级期间旧任务能上报重算进度但无法上报类型化状态此时status可能为NULL但progress字段仍然可用。RECALCULATE STREAM成功返回 ≠ 重算执行完成成功响应只代表请求被接受重算在后台运行。未完成请求的持久化未完成的重算请求在服务或流任务重启、重新部署后会恢复瞬态执行失败会自动重试。用recalc_id追踪请求在该请求处于Pending或Running期间避免重复提交相同请求。时间字段request_time是 mnode 接受请求的时间message在可用时包含状态或错误文本。终止记录的保留策略终止态Finished/Failed的重算记录从 mnode 首次观察到终止状态起保留1 小时每个流最多保留100 条终止记录Pending和Running记录不计入该上限。所有记录仅保存在内存中进程重启后可能消失。这一保留策略只针对终止记录与未完成请求的持久化是相互独立的。源码中 streamRecalcTracker.c 定义的STREAM_RECALC_MAX_TERMINAL_JOBS 100正是每流最多 100 条终止记录这一上限的实现依据。六、正确理解NULL与零值语义一分钟级指标统计的是最近 60 个完整秒不包含当前秒。任务启动、重启或重新部署后必须先产出一个完整窗口在此之前对应指标为NULL。具体规则如下完整窗口内无输入或无输出时对应速率为0完整窗口内 Runner 未形成任何结果样本时结果延迟为NULL流级输出速率要求每一个最终结果 Runner 都有有效的完整窗口任一必需 Runner 未就绪该字段即为NULL流级结果延迟同样要求至少有一个结果样本否则为NULL某指标不适用于某类任务时该字段为NULLNULL字段不会使无关字段失效例如没有结果延迟样本时输入速率仍然可用短暂的心跳中断期间管理节点可能保留最后一次成功的指标快照此时应结合任务级status与last_update判断新鲜度滚动升级期间尚未升级的任务无法提供新指标对应列可能暂时为NULL。七、推荐排查流程查询ins_streams中的status与message确认流是否健康检查realtime_lag_ms是否持续增长。若流使用外部数据源NULL是预期行为检查输入、输出与结果延迟。输入非零而输出为零不一定代表故障——窗口未闭合前、或过滤后无结果时都会出现这种情况查询ins_stream_tasks用status、last_update、node_id和任务级指标定位受影响的任务对于历史数据处理或手动重算问题分别检查history_progress_pct与ins_stream_recalculates。这套流程对应的正是本文开头提到的流级总览 → 任务级下钻 → 重算专项三级定位思路可与 01-meta.md 中的完整列定义相互对照进一步核对各字段的数据类型与取值范围。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考