ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

TDengine 流式计算运维与限制:高可用、权限控制、重算机制与边界条件实践指南

TDengine 流式计算运维与限制:高可用、权限控制、重算机制与边界条件实践指南 TDengine 流式计算运维与限制高可用、权限控制、重算机制与边界条件实践指南【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengineTDengine 的流式计算引擎采用触发与计算解耦的架构将流式任务全部交由 snode 执行。本文围绕流式计算上线后最关键的运维主题展开如何部署 snode 并实现高可用、如何做数据库级权限控制、如何用 WATERMARK 与重算机制应对乱序/更新/删除数据、哪些数据库与表操作会影响运行中的流任务以及完整的配置参数、规则限制与旧版本升级兼容步骤。读完本文你将掌握 TDengine 流式计算在生产环境中的部署规范、故障恢复手段与边界判断能力并能结合仓库中的源码与测试用例进一步验证每个行为。本文对应的官方文档为 docs/en/07-stream-processing/02-instructions.md是其全部内容的完整继承与源码级扩充语法层面的 CREATE STREAM 细节可参考 Stream Syntax可观测指标与排障流程见 Observability and Troubleshooting。一、流式计算的高可用架构1.1 存储与计算分离snodes 承担全部流计算TDengine 流式计算在架构上实现了计算与存储分离这要求集群中至少部署一个 snode。除数据读取外流式处理的全部功能都只在 snode 上运行。核心概念如下snode负责执行流式计算任务的节点。一个集群可以部署一个或多个 snode至少一个每个 dnode 上最多托管一个 snode每个 snode 内部拥有多个执行线程。部署位置snode 可以与 vnode/mnode 等其他节点类型部署在同一个 dnode 上但为了更好的资源隔离强烈建议将 snode 部署在独立 dnode 上使流式计算与写入、查询等操作互不干扰。高可用要求为保证流式计算的高可用建议在集群中跨不同物理节点部署多个 snode流任务会跨多个 snode 负载均衡每对 snode 互为副本保存流的状态与进度信息stream state and progress如果集群中只部署了单个 snode则系统无法保证高可用。从仓库测试用例可以印证上述行为。在 test/cases/18-StreamProcessing/01-Snode/ 目录下test_snode_mgmt_basic.py、test_snode_mgmt_replicas.py、test_snode_mgmt_replica3.py分别覆盖了 snode 管理、双副本与三副本场景而 test_stream_no_snode.py 则验证了未部署 snode 时无法创建流这一约束。部署与容量设计建议还可参考 Deployment and Design在定义流之前预先创建多个 snode 可获得更好的负载均衡流负载较重时通过横向增加 snode 扩容。1.2 部署、查看与删除 snode创建流任务之前必须先部署 snode语法如下CREATE SNODE ON DNODE dnode_id;查看集群中的 snodeSHOW SNODES;需要更详细信息时查询系统库视图SELECT * FROM information_schema.ins_snodes;删除 snode 时有明确的约束删除时 snode 及其副本必须同时在线以便同步流状态信息若两者任一离线删除操作会失败。DROP SNODE ON DNODE dnode_id;二、流式计算的权限控制流式计算的权限控制只与数据库级权限绑定。由于一个流可能关联多个数据库创建、删除、启停及手动重算流时的权限要求如下表关联数据库数量认证动作所需权限定义流所在的数据库1创建、删除、停止、启动、手动重算写Write触发表所在数据库1创建读Read输出表所在数据库1创建写Write计算源所在数据库1 个或多个创建读Read也就是说操作流本身包括重算要求对流所在数据库有 Write 权限而创建流时涉及的触发表与计算源只需 Read 权限输出表则需要 Write 权限。仓库中的权限测试覆盖了这一规则例如 test_snode_privileges_systable.py、test_snode_privileges_twodb.py 与 test_snode_privileges_recalc.py 分别验证了跨系统表、跨库及重算场景下的权限校验逻辑。三、重算Recomputation乱序、更新与删除的正确性保障3.1 为什么需要重算TDengine 的多数 TSDB 窗口类型与主键列关联。例如事件窗口依赖按主键有序的数据来决定窗口何时打开、何时关闭。使用基于窗口的触发时触发表数据应尽量有序写入这是流式计算最高效的写入方式。若写入乱序数据可能影响已触发窗口的结果正确性同理更新与删除操作也可能破坏结果正确性。为缓解乱序、更新、删除带来的问题TDengine 提供 WATERMARKWATERMARK是一个基于事件时间的、由用户定义的时长它表示系统在流处理中的进度反映用户对乱序数据的容忍度。当前 watermark 定义为最新已处理事件时间 – WATERMARK 间隔。只有事件时间早于当前 watermark 的数据才参与触发判定同样只有时间边界早于当前 watermark 的窗口或触发条件才会被触发。注意WATERMARK不适用于 PERIOD定时触发。在 PERIOD 模式下不进行任何重算。对于超出 WATERMARK 容忍范围的乱序、更新或删除场景使用重算保证结果正确性。重算的含义是对受乱序、更新、删除记录影响的数据区间重新触发并重新执行计算。已写入输出表的结果不会被删除而是再次写入新的结果。要让这种方式有效用户必须保证计算语句与源表不依赖处理时间——即同一触发即使执行多次也应产生有效结果。重算分为自动与手动两类。如果不需要自动重算可通过配置选项将其关闭。3.2 手动重算手动重算必须由用户显式发起需要时可用 SQL 命令启动RECALCULATE STREAM [db_name.]stream_name FROM start_time [TO end_time];使用要点如下可指定按事件时间计算的重算区间。若未指定结束时间end_time重算区间从给定开始时间start_time延伸到发起手动重算时流的当前处理进度。定时触发PERIOD不支持手动重算其余触发类型均支持。对计数窗口触发开始时间与结束时间都必须指定重算仅作用于流已经处理过的区间若指定区间包含流尚未开始处理的区间这些部分会被自动忽略。重算期间触发窗口会在指定区间内重新划分可能导致新窗口与之前计算的窗口错位用户可能需要手动删除输出表中该区间的旧结果以避免重复同样若未指定结束时间重算请求会被忽略。对于从某个开始时间重算且没有明确结束的场景推荐做法是删除流、重新创建并指定FILL_HISTORY_FIRST。SQL 返回成功只代表重算请求已被接受不代表重算已完成。应使用recalc_id与information_schema.ins_stream_recalculates监控请求状态。重算在后台运行。若请求到达终态前服务或流任务重启、重新部署未完成的请求会被恢复并继续处理。瞬时执行失败会自动重试。不要因为请求仍处于Pending或Running就重复提交同一请求应先检查status、progress与message字段。关于手动重算的请求生命周期、状态枚举Pending/Running/Finished/Failed以及终态记录保留策略终态记录自 mnode 观察到终态起保留 1 小时每个流最多 100 条且仅驻留内存详见 Observability and Troubleshooting 中的View Manual Recalculations一节。3.3 重算命令在解析层的实现从源码结构看RECALCULATE STREAM是一条独立的 SQL 语句在 sql.y 的语法定义与 parAstParser.c 的 AST 解析中对应QUERY_NODE_RECALCULATE_STREAM_STMT节点最终生成带时间范围的重算请求下发给 mnode 执行。这与文档中SQL 返回成功仅代表请求被接受、由 mnode 统一管理重算进度的描述一致——重算请求是异步的由管理节点跟踪其生命周期。四、非典型数据写入场景的处理4.1 乱序数据Out-of-Order Data乱序数据指以非顺序方式写入触发表的记录。计算本身不依赖源表是否有序但用户必须根据业务需求确保源表数据在触发前完整写入。乱序数据的影响与处理方式因触发类型而异触发类型影响与处理方式定时触发滑动触发滑动步长不为 1 的计数窗口触发忽略不进行任何处理滑动步长为 1 的计数窗口触发如COUNT_WINDOW(1)、COUNT_WINDOW(n, 1)默认通过重算处理可选使用STREAM_OPTIONS(IGNORE_DISORDER)忽略其他窗口触发默认通过重算处理可选忽略不进行任何处理4.2 数据更新Data Updates数据更新指同一时间戳被多次写入、其他列的值可能变化的情况。更新只影响触发表与触发行为不直接影响计算过程本身触发类型影响与处理方式定时触发滑动触发滑动步长不为 1 的计数窗口触发忽略不进行任何处理滑动步长为 1 的计数窗口触发如COUNT_WINDOW(1)、COUNT_WINDOW(n, 1)默认视为乱序数据通过重算处理可选使用STREAM_OPTIONS(IGNORE_DISORDER)忽略其他窗口触发视为乱序数据通过重算处理4.3 数据删除Data Deletions数据删除同样只影响触发表与触发行为不直接影响计算过程触发类型影响与处理方式定时触发滑动触发滑动步长不为 1 的计数窗口触发忽略不进行任何处理滑动步长为 1 的计数窗口触发如COUNT_WINDOW(1)、COUNT_WINDOW(n, 1)默认忽略不进行任何处理可选使用STREAM_OPTIONS(DELETE_RECALC)视为乱序数据并通过重算处理其他窗口触发默认忽略不进行任何处理可选视为乱序数据并通过重算处理4.4 过期数据Expired Dataexpired_time设置定义了一个数据过期区间。对每个由流触发生成的组系统通过比较最新数据的事件时间与过期阈值来判断新数据是否过期。阈值计算方式为最新事件时间 – expired_time早于该阈值的数据一律视为过期数据。关键性质如下过期数据只针对触发表中的实时数据。历史数据与其他表的数据不存在过期概念。过期与否在流触发时评估取决于数据写入的时机按事件时间顺序写入的数据永远不会过期只有乱序数据才可能被判为过期。过期数据不会自动触发新的计算或重算。即所有触发类型下过期数据都被忽略不计算、不重算。如果不需要从计算或重算中排除任何时间区间就无需指定 expired_time若已定义过期数据但仍希望对其中一部分进行计算或重算可以使用手动重算。过期数据只影响是否自动触发不影响计算区间本身。因此若某次触发的计算区间包含触发表中的过期数据这些数据仍会参与计算。五、数据库与表操作对运行中流的影响流创建之后用户可能对与流关联的数据库和表执行各种操作。这些操作对流的影响及流的处理方式汇总如下操作操作影响与流处理方式在触发超级表非虚拟下新建子表并写入数据新子表自动纳入当前流处理可加入已有组或创建新组在虚拟触发超级表下新建子表并写入数据忽略不做额外处理删除触发超级表的子表默认忽略可选部分触发类型可配置为自动重算或删除对应结果表仅适用于按子表分组的流对ROLLUP BY流已生成的输出子表会保留删除触发表忽略不做额外处理向触发表添加列忽略不做额外处理删除触发表某列忽略不做额外处理修改触发超级表下子表的标签值若该标签列被流用作分组键操作不允许并报错否则忽略修改触发表列的 schema忽略不做额外处理读取时检测到 schema 不匹配会报错修改或删除源表忽略不做额外处理修改或删除输出表忽略不做额外处理写入时检测到 schema 不匹配会报错表不存在时会自动重建拆分 vnode若 vnode 所在数据库是源数据库或触发表数据库不允许若触发或计算使用虚拟表不允许确认无影响后可用SPLIT VGROUP N FORCE强制执行删除数据库若被删数据库是流的源数据库或是不与流所在数据库相同的触发表数据库不允许若流涉及来自非目标数据库的虚拟表触发或计算不允许确认无影响后可用DROP DATABASE name FORCE强制执行除上表明确受限或特殊处理的操作外其余操作以及表中标记为忽略不做额外处理的操作均不受限制。但若此类操作可能影响流计算则由用户自行决定处理方式忽略影响或执行手动重算以恢复正确性。六、流式计算配置参数流式计算相关的配置参数如下表完整细节见 taosd。从仓库中参数注册的源码 tglobal.c 可以看到每个参数的实际默认值与取值范围参数文档说明源码中的默认值与范围见 tglobal.cnumOfMnodeStreamMgmtThreadsmnode 上的流管理线程数默认 2范围 [2, 5]numOfStreamMgmtThreadsvnode/snode 上的流管理线程数默认 2范围 [2, 5]numOfVnodeStreamReaderThreadsvnode 上的流读取线程数默认 2范围 [2, INT32_MAX]numOfStreamTriggerThreads流触发线程数默认 4范围 [4, INT32_MAX]numOfStreamRunnerThreads流执行线程数默认 4范围 [4, INT32_MAX]streamBufferSize流处理可用最大缓冲区仅用于缓存%%trows的结果单位MB默认 128见 tglobal.c范围 [128, INT32_MAX]streamNotifyMessageSize控制事件通知消息的大小默认 8范围 [8, 1024 × 1024]见 tglobal.cstreamNotifyFrameSize控制发送事件通知消息时底层 frame 的大小默认 8范围 [8, 1024 × 1024]见 tglobal.c结合部署实践理解这些参数详见 Deployment and Design线程数应按部署方式与负载调节线程数越高 CPU 资源消耗越大越低则消耗越小。最大流缓冲区大小也应按部署方式配置负载较重或并发流很多时应增大缓冲区——尤其是 snode 部署在独立节点时。仓库的测试用例对参数范围与动态调整有专门覆盖例如 test_snode_params_check_minvalue.py、test_snode_params_check_maxvalue.py、test_snode_params_check_default.py、test_snode_params_alter_value.py 以及缓冲区专项 snode_params_buffersize.py。另需注意部分参数支持动态生效CFG_DYN_SERVER_LAZY或CFG_DYN_SERVER级别修改配置的生效范围可参考 Config Scope。七、规则与限制7.1 一般规则以下规则与限制适用于流式计算创建流之前集群必须至少部署一个 snode且创建时刻必须有一个可用运行中的 snode。每个流属于特定数据库。因此创建流之前数据库必须已存在且同一数据库内流不能重名。流的触发表与源表可以相同或不同也可以属于不同数据库。流的输出表可以位于与流、触发表或源表不同的数据库但不能与触发表或源表相同。输出表无论是超级表还是普通表在流创建时自动创建。若希望写入已存在的表其 schema 必须完全匹配。每个组的输出子表无需提前创建计算写出结果时会自动创建。每个触发组的计算结果写入同一个子表若未指定触发组所有结果写入单个普通表。若不同组被配置为生成同名的子表其结果会写入同一个子表。用户必须确认这是预期行为否则应确保每个组生成唯一命名的子表。除指定子表名外用户还可以定义输出超级表的标签列以及每个子表的标签值。流式计算支持嵌套即可基于已有流的输出表创建新流。对于滑动步长不为 1 的计数窗口触发乱序数据、更新与删除均被忽略滑动步长为 1如COUNT_WINDOW(1)、COUNT_WINDOW(n, 1)时乱序数据与更新默认触发重算可用IGNORE_DISORDER忽略删除默认忽略可用DELETE_RECALC触发重算。在非FILL_HISTORY_FIRST模式下历史窗口与实时窗口可能不对齐。对于超级表窗口触发仅 interval 与 session 窗口支持按标签、按 rollup 标签、按子表分组或不分组其他窗口类型只支持按子表分组。ROLLUP BY与PARTITION BY、DELETE_OUTPUT_TABLE互斥。查询中不支持伪列 qstart、qend、qduration。7.2 临时限制以下为当前版本尚未支持的能力尚不支持按普通数据列分组。尚不支持 Geometry 数据类型。NOTIFY_OPTIONS中的ON_FAILURE_PAUSE选项尚不支持。八、兼容性说明从旧版本升级与v3.3.6.0相比流式计算已被完全重新设计。从旧版本升级前必须按顺序执行以下步骤并在新版流式计算下重新创建流删除所有已有的流式计算任务。删除所有 TSMA。删除所有 snode。移除 snode 相关目录dataDir 配置路径下的 snode 目录默认/var/lib/taos/snode旧checkpointBackupDir配置项指定的目录默认/var/lib/taos/backup/checkpoint/。删除所有结果表。注意若未执行上述步骤taosd 将无法启动。升级后的新引擎提供了更大的灵活性与更高的可用性同时也对正确使用提出了更高要求建议升级后结合 Deployment and Design 中的设计检查清单重新审视流定义并通过 Observability and Troubleshooting 中的系统视图持续监控流运行状态。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进