资讯详情

TDengine 数据订阅引擎内部原理:Topic、Consumer Group、WAL 与 Rebalance 机制解析

📅 2026/9/21 16:20:38 | 华诺云谱 👁 阅读
TDengine 数据订阅引擎内部原理:Topic、Consumer Group、WAL 与 Rebalance 机制解析
数据库时序数据库物联网大数据实时分析云原生【免费下载链接】tdengineTDengine is an open source, high-performance, cloud native time-series database optimized for Internet of Things (IoT), Connected Cars, Industrial IoT and DevOps.项目地址https://gitcode.com/taosdata/tdengine点击查看免费下载TDengine 将时序数据库与消息队列能力合二为一其内置的数据订阅Data Subscription / TMQ引擎允许用户像使用 Kafka 一样订阅数据库、超级表或 SQL 查询结果。本文以 docs/en/15-internals/05-topic.md 为骨架结合仓库内订阅实现源码clientTmq.h、mndConsumer.c、mndSubscribe.c与官方 Topic 语法文档docs/en/06-data-subscription/01-topic.md深入剖析 Topic、Producer、Consumer、Consumer Group、消费进度、客户端/服务端架构、Rebalance 过程与基于 WAL 的数据消费原理。读完本文你将理解 TMQ 的端到端工作流程掌握消费进度提交与earliest/latest语义并能对照源码理解 vnode 分配、心跳与状态机等底层实现。基本概念Topic与 Kafka 类似使用 TDengine 数据订阅需要先定义一个Topic。TDengine 的 Topic 可以是数据库、超级表supertable或一条查询语句subquery数据库订阅与超级表订阅主要用于数据迁移场景可以在另一个集群中完整恢复整个数据库或超级表查询语句订阅是 TDengine 数据订阅的亮点灵活性更高。因为数据过滤与预处理由 TDengine 完成而非应用层可以有效减少传输的数据量和应用复杂度。Topic 的数据表分布在多个 vnode 上vnode 对应 Kafka 中的 partition每个 vnode 的数据按顺序写入 WAL 文件。由于 WAL 文件中不仅存储数据还存储元数据、写入消息等因此数据的版本号version并不连续。TDengine 会自动为 WAL 文件建立索引以支持快速随机访问通过灵活可配置的文件切换roll与保留retention机制用户可以按需指定 WAL 文件的保留时间与大小使 WAL 成为一个保留事件顺序的持久化存储引擎。仓库配置文档 docs/en/12-operations-and-tooling/03-components/01-taosd.md 中即包含walRetentionPeriodWAL 保留时长秒与walRetentionSizeWAL 保留大小字节等参数只有该保留窗口内的增量数据才能被 TMQ/taosX 订阅消费参见 docs/en/11-security-guide/03-full-trace-reliability.md。对于查询语句订阅消费时 TDengine 根据当前消费进度直接从 WAL 文件中读取数据通过统一查询引擎执行过滤、变换等操作再推送给消费者。Producer**Producer生产者**是与订阅 Topic 数据表相关联的数据写入应用。生产者可以通过多种方式生成数据并写入数据表所在 vnode 的 WAL 文件这些方式包括SQL 写入Stmt参数化/批量写入Schemaless 写入CSV 导入流式计算stream computing结果写入ConsumerConsumer消费者负责从 Topic 中获取数据。订阅 Topic 后消费者可以消费分配给该消费者的所有 vnode 上的数据。为实现高效有序的数据获取消费者采用推送push 拉取poll相结合的方式当 vnode 中有大量未消费数据时消费者会按顺序向 vnode 发送 push 请求一次性拉取大批量数据同时消费者在本地记录每个 vnode 的消费位置确保所有数据按序推送当 vnode 中无数据可消费时消费者进入等待状态。一旦有新的数据写入 vnode系统会立即通过 push 方式将数据推送给消费者保证数据的及时性。从客户端源码可见拉取空闲时使用EMPTY_BLOCK_POLL_IDLE_DURATION100ms控制空块轮询间隔而DEFAULT_ASKEP_INTERVAL1000ms则控制客户端向服务端询问/获取信息的周期见 source/client/inc/clientTmq.h。Consumer Group创建消费者时必须指定一个消费组consumer group。同一消费组内的消费者共享消费进度从而保证数据在消费者之间均匀分布。如前所述一个 Topic 的数据分布在多个 vnode 上为了提升消费速度、实现多线程分布式消费可以在同一消费组中增加多个消费者这些消费者会先均分 vnode再消费分配给自己的 vnode。例如数据分布在 4 个 vnode 上2 个消费者时每个消费者消费 2 个 vnode3 个消费者时2 个消费者各消费 1 个 vnode剩余 1 个消费者消费剩下的 2 个 vnode5 个消费者时4 个消费者各分得 1 个 vnode剩余 1 个消费者不参与消费。向消费组新增消费者后系统会通过rebalance 机制自动重新分配消费者该过程对用户透明、无需人工干预。此外一个消费者可以订阅多个 Topic 以满足不同场景的数据处理需求即使在崩溃、重启等复杂环境下TDengine 数据订阅仍能保证**至少一次at least once**消费确保数据完整可靠。消费进度消费组在 vnode 中记录消费进度以便在消费者重启或故障恢复时准确恢复消费位置。消费过程中消费者可以提交消费进度即 vnode 上 WAL 的版本号对应 Kafka 的 offset。消费进度提交可以手动进行也可以通过参数设置为周期性自动提交。消费者首次消费时可通过订阅参数决定消费位置即消费最新数据还是最旧数据。对于同一 Topic 与任意消费组每个 vnode 的消费进度是唯一的。因此当某个 vnode 上的消费者提交进度并退出后同组其他消费者将从该进度继续消费若前一消费者未提交进度新消费者将根据订阅参数设置决定起始消费位置。需要注意的是不同消费组即使消费同一个 Topic 也不共享消费进度这一设计保证了各消费组的独立性使其可以互不干扰地独立处理数据。数据订阅架构数据订阅系统在逻辑上分为**客户端client与服务端server**两个核心模块客户端负责创建消费者、获取这些消费者独占的 vnode 列表、从服务端拉取所需数据并维护必要的状态信息服务端专注于管理与 Topic、消费者相关的信息处理客户端的订阅请求实现 rebalance 机制以动态分配消费者节点保证消费过程的连续性与数据一致性同时跟踪和管理消费进度。客户端与服务端成功建立连接后用户必须先指定消费组与 Topic 来创建对应的消费者实例然后客户端向服务端提交订阅请求。此时消费者的状态被标记为rebalancing处于 rebalance 阶段。消费者会周期性向服务端发送请求以获取待消费的 vnode 列表直到服务端完成 vnode 分配。分配完成后消费者状态更新为ready表示订阅过程成功完成客户端可以正式向 vnode 发送数据消费请求。在数据消费过程中消费者不断向每个分配的 vnode 发送请求以获取新数据收到数据并消费完成后继续向该 vnode 发送请求以保持持续消费。若在预设时间内未收到数据消费者会在 vnode 上注册一个消费句柄handle一旦 vnode 产生新数据便立即推送给消费者从而保证消费的即时性并有效降低消费者频繁主动拉取带来的性能损耗。客户端拉取数据的模式本质上是 pull 与 push 的高效结合。消费者收到数据时还会收到数据的版本号并将其记录为各 vnode 的当前消费进度该进度保存在消费者内存中、仅对当前消费者有效。若消费者需要退出并希望在之后恢复上次消费进度则必须在退出前向服务端提交消费进度commit 操作使进度持久化存储在服务端支持自动与手动两种提交方式。此外客户端实现了心跳保活heartbeat keep-alive机制通过定期向服务端发送心跳证明消费者在线。若服务端在一段时间内未收到某消费者的心跳将判定其离线对于长时间不拉取数据的消费者时长可由参数控制服务端也会将其标记为离线并从消费组中移除。服务端依靠心跳机制监控所有消费者状态从而有效管理整个消费组。客户端侧默认心跳间隔为DEFAULT_HEARTBEAT_INTERVAL3000ms见 source/client/inc/clientTmq.h心跳超时判定则由session.timeout.ms参数控制默认 12000ms。从服务端模块划分看mnode主要处理订阅过程中的控制消息包括创建/删除 Topic、订阅消息、查询 endpoint 消息、心跳消息等vnode专注于处理消费消息consumption message与提交消息commit message。当 mnode 收到消费者的订阅消息时若该消费者此前未订阅过或已订阅但订阅的 Topic 发生变化其状态都会被置为 rebalancing随后 mnode 会对处于 rebalancing 状态的消费者执行 rebalance 操作心跳超时超过固定时间的消费者或被主动关闭的消费者将被删除。相关状态机定义可见 source/dnode/mnode/impl/inc/mndConsumer.hMQ_CONSUMER_STATUS_REBALANCE、MQ_CONSUMER_STATUS_READY。消费者定期向 mnode 发送查询 endpoint 消息以获取 rebalance 后的最新 vnode 分配结果同时定期发送心跳消息告知其在线状态消费者的部分信息也会随心跳上报至 mnode用户可以在 mnode 上查询这些信息以监控各消费者状态便于有效管理与监控。Rebalance 过程每个 Topic 的数据可能分散在多个 vnode 上通过执行rebalance 过程服务端将这些 vnode 合理地分配给各个消费者保证数据均匀分布与高效消费。以下图为例c1 代表消费者 1c2 代表消费者 2g1 代表消费组 1。初始时 g1 中只有 c1 消费数据c1 向 mnode 发送订阅信息mnode 将包含数据的全部 4 个 vnode 分配给 c1。当 c2 加入 g1 后c2 向 mnode 发送订阅信息mnode 检测到 g1 需要重新分配并发起 rebalance 过程随后将其中 2 个 vnode 分配给 c2 消费分配信息也由 mnode 发送给 vnodec1 与 c2 分别从各自被分配的 vnode 开始消费。Rebalance 定时器每 2 秒检查一次是否需要重新分配。在 rebalance 过程中若消费者状态不是 ready则无法消费只有 rebalance 正常结束且消费者获取到被分配 vnode 的 offset 后才能正常消费否则消费者会重试指定次数后报错。客户端源码中SUBSCRIBE_RETRY_MAX_COUNT240与SUBSCRIBE_RETRY_INTERVAL500ms即对应订阅/重平衡失败时的重试策略见 source/client/inc/clientTmq.h。rebalance 的实际分配逻辑按 vgroup 平均分配、vnode 分裂触发重平衡等可在 mndSubscribe.c 中看到相关实现与日志例如 mq rebalance add new consumer、mq rebalance vgId ... moved from consumer ... to consumer ... 等。消费者状态处理消费者的状态转换过程如下图所示。刚完成订阅的消费者处于rebalancing状态表示尚未准备好消费数据一旦 mnode 检测到处于 rebalancing 状态的消费者就会发起 rebalance 过程。rebalance 成功后消费者状态变为ready。随后消费者周期性查询 endpoint 消息以获取 ready 状态与被分配的 vnode 列表即可正式开始消费数据。若消费者心跳丢失超过 12 秒在 rebalance 过程后其状态将被更新为clear随后被系统删除当消费者主动退出时会发送 unsubscribe 消息该消息会清除该消费者订阅的所有 Topic 并将其状态置为 rebalancing。随后系统检测到 rebalancing 状态消费者并启动 rebalance 过程成功后该消费者状态更新为 clear最终被系统删除。这一系列措施保证了消费者的有序退出与系统稳定性。心跳超时阈值与session.timeout.ms参数默认 12000范围 [6000, 1800000]对应长时间不 poll 的判定则由max.poll.interval.ms参数默认 300000控制。消费数据时序数据存储在 vnode 上消费的本质就是读取 vnode 上 WAL 文件中的数据。WAL 文件扮演消息队列的角色消费者记录 WAL 数据的版本号本质上就是追踪消费进度。WAL 文件中的数据包括数据data与元数据meta如表创建、修改操作。订阅根据 Topic 的类型与参数获取相应数据若订阅涉及带过滤条件的查询订阅逻辑会通过通用查询引擎过滤掉不满足条件的数据。如下所示vnode 可以通过设置参数自动提交消费进度也可以由消费者在确认数据处理完成后手动提交。若消费进度存储在 vnode 中则同一消费组中不同消费者切换时会延续之前的进度否则根据配置参数消费者可以选择消费最旧数据或最新数据。earliest参数表示消费者从 WAL 文件中最旧的数据开始消费latest参数表示从最新数据即新写入的数据开始消费。这两个参数只在消费者首次消费或未提交消费进度时生效。若消费过程中提交了消费进度——例如消费完 WAL 中第 3 条数据后提交一次commit offset3——那么下次在同一 vnode 上、同一消费组与 Topic 的新消费者将从第 4 条数据开始消费。该设计既允许消费者按需灵活选择消费起点又保持了消费进度的持久性与消费者之间的同步。实操Topic 语法与消费者参数速查结合 docs/en/06-data-subscription/01-topic.md 与 docs/en/06-data-subscription/02-native.md将上述内部原理落地为可直接执行的 SQL 与参数配置创建查询 TopicCREATE TOPIC [IF NOT EXISTS] topic_name AS subquery;例如订阅power库meters表中电压大于 200 的数据行仅返回时间戳、电流、电压列CREATE TOPIC power_topic AS SELECT ts, current, voltage FROM power.meters WHERE voltage 200;查询 Topic 支持过滤条件与标量函数但不支持聚合函数、时间窗口聚合以及DISTINCT、GROUP BY、ORDER BY、PARTITION BY、LIMIT/SLIMIT等子句。注意订阅所引用或参与计算的列不能被删除或修改从 v3.4.0.0 起可修改/删除/新增但需执行RELOAD TOPIC生效虚拟表不支持查询订阅删除子查询中的表后订阅数据为空重建同名表后仍为空表 ID 已变化需RELOAD TOPIC重新订阅。创建超级表 / 数据库 TopicCREATE TOPIC [IF NOT EXISTS] topic_name [WITH META | ONLY META] AS STABLE stb_name [where_condition]; CREATE TOPIC [IF NOT EXISTS] topic_name [WITH META | ONLY META] AS DATABASE db_name;WITH META额外返回建超级表及其子表的语句主要用于 taosX 做超级表/数据库迁移ONLY META只订阅元数据变更不传输时序数据超级表订阅的WHERE子句只能使用标签或tbname过滤子表不能使用普通列超级表/数据库订阅是高级模式、更易出错如需使用建议咨询技术支持。删除、查看与重载 TopicDROP TOPIC [IF EXISTS] [FORCE] topic_name; -- FORCE 支持强删正在被订阅的 Topicv3.3.6.0 SHOW TOPICS; RELOAD TOPIC [IF EXISTS] topic_name AS subquery; -- v3.4.0.0仅查询 Topic实例中可创建的 Topic 总数由tmqMaxTopicNum控制范围 1–10000默认 20见 docs/en/12-operations-and-tooling/03-components/01-taosd.md。创建消费者时的关键参数详细清单见 docs/en/10-developer-guide/07-subscription-api.md参数说明默认值td.connect.ip/td.connect.port服务端 FQDN 与端口—td.connect.user/td.connect.pass/td.connect.token用户名/密码/令牌认证—group.id消费组 ID同组共享消费进度必填auto.offset.reset消费组订阅初始位置earliest/latest/nonelatestv3.2.0.0enable.auto.commit是否自动提交消费进度trueauto.commit.interval.ms自动提交间隔毫秒5000session.timeout.ms心跳丢失判定超时毫秒12000max.poll.interval.ms消费者两次 poll 最大间隔毫秒300000fetch.max.wait.ms单次 fetch 服务端最大等待时间毫秒1000min.poll.rows单次服务端返回的最小行数4096enable.replay是否启用数据重放replay关闭Replay数据重放TDengine 订阅支持按原始写入时间间隔重新推送消息基于 WAL 实现。例如三行数据写入时间分别为00:00:00.000、00:00:05.000、00:00:08.000replay 会立即返回第一条约 5 秒后返回第二条再过约 3 秒返回第三条。仅查询 Topic 支持 replayreplay 进度不保存。此外还可用taosshell 的subscribe topic -g group_id命令快速验证 Topic 是否能产出数据详见 docs/en/12-operations-and-tooling/04-tools/01-taos-cli.md 的数据订阅小节。总结TDengine 数据订阅引擎以WAL 作为持久化消息队列、以vnode 作为分区partition、以版本号作为 offset通过 mnode 统一管理 Topic 与消费者控制消息、vnode 处理数据与提交消息借助rebalance 机制实现消费组内的动态负载均衡并以**心跳保活 状态机rebalancing → ready → clear**保证消费者的在线检测与有序退出。理解这些内部原理可以帮助你在数据迁移、实时流处理、应用解耦等场景中更合理地设计 Topic 与消费组并准确利用earliest/latest、手动/自动提交等机制控制消费行为。赞分享数据库时序数据库物联网大数据实时分析云原生【免费下载链接】tdengineTDengine is an open source, high-performance, cloud native time-series database optimized for Internet of Things (IoT), Connected Cars, Industrial IoT and DevOps.项目地址https://gitcode.com/taosdata/tdengine点击查看免费下载相关推荐TDengine 数据订阅引擎内幕TMQ Topic、消费者组与 Rebalance 机制TDengine 数据订阅引擎内幕TMQ Topic、消费者组与 Rebalance 机制 TDengine 的数据订阅TMQTDengine Messa数据库时序数据库大数据物联网云原生TDengine 数据订阅TMQ内部原理主题、消费者组、WAL 消费与 Rebalance 机制深度解析TDengine 数据订阅TMQ内部原理主题、消费者组、WAL 消费与 Rebalance 机制深度解析 数据订阅TDengine Message Qu数据库时序数据库大数据物联网云原生TDengine 数据订阅主题Topic完全指南语法、管理与 WAL 回放机制TDengine 数据订阅主题Topic完全指南语法、管理与 WAL 回放机制 本篇围绕 TDengine 数据订阅体系中的主题Topic展开从数据库时序数据库大数据物联网云原生上一篇终极指南Neovim-from-scratch团队协作开发配置共享的10个技巧下一篇RemoteCam免费开源的Android摄像头桌面串流工具让手机秒变OBS摄像头与虚拟 webcam创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。