工业实时数据分析的架构正在经历一次范式转变。传统架构中,数据从采集到分析需要穿越四个环节:传感器采集→消息队列→流计算引擎→数据库。每个环节增加延迟、引入运维复杂度、扩大故障域。时序数据库内置流计算能力的出现,使数据在写入存储时即触发计算,省去了数据从存储到流计算引擎的往返传输。架构层数从四级压缩到两级——采集→内置流计算的时序数据库——延迟从分钟级降到秒级。
传统实时分析架构的链路问题
传统工业实时分析架构采用"职责分离"设计:采集网关负责协议解析和数据标准化,Kafka/RabbitMQ负责数据缓冲和削峰,Flink/Spark Streaming负责流式计算,数据库负责持久化存储。每一层各司其职,看似清晰,实则链路过长。
延迟累积是首要问题。数据从传感器到数据库的端到端延迟由各环节延迟叠加:采集网关解析延迟10-50毫秒,Kafka生产+消费延迟50-200毫秒,流计算引擎处理延迟100-500毫秒,数据库写入延迟5-20毫秒。总延迟165-770毫秒。如果采集频率10Hz(100毫秒/次),端到端延迟可能达到5-8个数据周期。告警从数据产生到触发,可能已经过了近1秒。对于缓慢变化的物理量(如温度),1秒延迟可接受;但对于故障状态量(如断路器跳闸),1秒延迟意味着告警落后于事件。
组件数量多导致运维复杂。一个完整的传统架构至少涉及采集网关、Kafka集群、ZooKeeper集群、Flink集群、时序数据库五个独立组件。每个组件需要独立的部署、监控、告警、容量规划和故障处理。Kafka集群的分区不均衡、Flink的checkpoint失败、ZooKeeper的脑裂——任何一个组件的故障都可能导致数据链路中断。
数据一致性在多组件间难以保证。Kafka消费偏移和数据库写入是两个独立操作。如果Flink从Kafka消费了数据但写入数据库前崩溃,数据在Kafka中已标记消费但未持久化——数据丢失。Flink的exactly-once语义可以缓解但不能完全消除这一问题,实现成本高且有性能开销。
资源浪费在数据搬运上。Kafka将数据从采集网关传到流计算引擎,流计算引擎处理后将结果传到数据库。两次网络传输和序列化/反序列化消耗CPU和带宽。对于10万测点×1Hz采样的系统,每秒10万条数据在组件间搬运两次,网络流量约20MB/s,序列化CPU开销约占单节点CPU的5-10%。
内置流计算的技术原理
时序数据库内置流计算的核心思路是:数据在写入存储引擎的同时,触发预定义的计算逻辑,计算结果直接写入结果表或推送给消费端。数据不需要离开数据库进程,省去了网络传输和序列化开销。
架构维度 | 传统架构(采集+MQ+流计算+DB) | 内置流计算(采集+DB) |
组件数量 | 4-5个 | 2个 |
端到端延迟 | 165-770ms | 10-50ms |
数据搬运次数 | 2次网络传输 | 0次(进程内) |
运维复杂度 | 高(多集群管理) | 低(单系统管理) |
故障域 | 大(任一组件故障影响全局) | 小(单系统内) |
exactly-once | 需复杂配置 | 写入即计算,天然保证 |
TDengine的内置流计算通过两种机制实现:连续查询(Continuous Query)和数据订阅(TMQ)。
连续查询是写入触发的预计算。配置一个连续查询:每5分钟计算各设备过去5分钟的温度均值和最大值,结果写入结果表。底层机制是定时器触发查询——每隔5分钟执行一次SQL,扫描过去5分钟的原始数据做聚合。连续查询的结果表本身也是普通超级表,可以被后续查询或二次连续查询消费。这种链式预计算——原始表→5分钟级结果表→1小时级结果表→1天级结果表——实现了多级聚合自动化。
数据订阅(TMQ)是写入触发的实时推送。消费者订阅某个超级表的数据变更,当新数据写入时,数据库将数据推送到消费者。推送机制基于数据库内部的写入日志(类似于binlog),保证数据不丢失——消费者断线重连后从上次消费位置继续。TMQ适用于需要毫秒级实时响应的场景——告警引擎订阅测点数据流,数据写入后立即被消费,阈值判断和告警触发在秒内完成。
典型流计算场景的工程实现
实时聚合。 每5分钟计算产线温度均值、设备振动RMS、产线OEE。实现方式:配置连续查询,每5分钟执行一次,结果写入5分钟级聚合表。查询实时看板时直接读聚合表,不需要实时扫描原始数据。1天的原始数据(1秒采样×10万测点=8.64亿条)被压缩为288×10万=2880万条聚合结果,查询效率提升30倍。
异常检测与阈值告警。 订阅关键测点数据流,实时判断是否越限。实现方式:消费者程序通过TMQ订阅超级表数据,每条数据到达后检查是否超过预设阈值。超过阈值时生成告警事件,写入告警表。告警延迟从传统架构的数百毫秒降到10-50毫秒。
复杂异常检测需要窗口计算。单个数据点越限是简单判断,但很多工业异常需要基于一段时间窗口的数据模式判断。例如:振动加速度连续10秒超过5g——滑动窗口判断;温度变化率1分钟内超过10°C/分钟——差分窗口判断;三相电流不平衡度5分钟内持续超过5%——持续窗口判断。
窗口计算。 TDengine支持两种窗口类型。滚动窗口(Tumbling Window)将时间轴划分为固定大小的非重叠窗口,每个窗口独立计算。INTERVAL(5m)就是5分钟滚动窗口。滑动窗口(Sliding Window)窗口按固定步长滑动,窗口之间有重叠。INTERVAL(5m) SLIDING(1m)每1分钟计算一次过去5分钟的数据。滑动窗口比滚动窗口的计算量大(数据被多次计算),但响应更灵敏——异常在第1分钟内就能被检测到,而非等到第5分钟窗口结束。
流计算与离线计算的工程边界
流计算和离线计算不是替代关系,而是互补关系。混淆两者的边界会导致架构设计不当。
流计算追求实时性,接受近似结果。连续查询在时间窗口结束时计算,窗口边界附近的数据可能不完整。例如5分钟窗口结束时,最后1秒的数据可能还在写入中,连续查询使用的是窗口开始到触发时刻的已写入数据。对于实时监控和快速响应场景,这种近似可接受——温度均值差0.1°C不影响告警判断。
离线计算追求精确性,接受高延迟。T+1的日报表需要完整覆盖当天所有数据,包括23:59:59写入的数据。离线批处理在第二天凌晨执行,确保数据完整性。对于财务结算、合规报告等精确性要求高的场景,离线计算不可替代。
工程实践中的分层策略:流计算负责实时监控和快速响应(秒级到分钟级延迟),离线计算负责精确报表和深度分析(小时级到天级延迟)。两者使用同一份原始数据,计算逻辑可以不同——流计算用近似算法(如在线均值),离线计算用精确算法(如全量均值)。
TDengine在这一分层架构中的定位:既提供流计算能力(连续查询+数据订阅),也提供离线计算能力(全量SQL查询)。数据不需要离开数据库——流计算在数据库内完成,离线计算也在数据库内完成。相比Flink+Hive的传统组合,组件数量少、数据不搬移、运维简单。
流计算的性能调优
内置流计算的性能取决于三个因素:数据写入速率、计算逻辑复杂度、结果写入开销。
数据写入速率决定流计算触发频率。1秒采样的数据,连续查询每5分钟触发一次,窗口内数据量600条/测点。如果测点10万个,单次连续查询扫描6000万条数据,聚合计算在单节点上可能需要2-5秒。如果5分钟内写入还在持续,下一周期的连续查询可能与上一个还未完成,形成积压。解决方案:降低连续查询频率(改为10分钟一次)或增加计算节点(TDengine集群模式下连续查询在多个Vnode上并行执行)。
计算逻辑复杂度影响CPU占用。简单的AVG/MAX/MIN计算CPU开销低,复杂的自定义函数(UDF)或窗口内的排序操作CPU开销高。流计算逻辑应尽量简化——复杂分析放到离线计算中做。
结果写入开销是容易被忽视的瓶颈。连续查询的结果写入结果表,如果结果表写入吞吐跟不上连续查询产出速度,结果表写入会成为瓶颈。实践中应监控结果表的写入延迟,确保连续查询的计算产出能被及时持久化。
结语
流计算从外置组件向时序数据库内置的迁移,是工业实时分析架构的演进方向。驱动力不是"少装一个组件"的运维简化,而是数据移动路径缩短带来的延迟降低和一致性保障。连续查询处理周期性聚合,数据订阅处理事件驱动场景,窗口计算处理模式检测——三种机制覆盖了工业流计算的主要需求。但内置流计算不意味着离线计算消失,两者在延迟性和精确性维度上互补。架构设计的关键是明确每种计算需求的时间约束——秒级响应用数据订阅,分钟级聚合用连续查询,天级精确报表用离线批处理。将正确的计算放到正确的时间维度上执行,是工业实时数据分析架构设计的核心原则。

























