TDengine 流计算让异常管控“真实时”:江苏正元数币在智慧校园的领先实践

正元数币

2026-08-27 /

小T导读:当一个智慧校园物联平台管理着上万个温控器和断路器、每秒钟都在产生新的温度读数和用电量数据时,真正的挑战不是”能不能存下”,而是”能不能在数据产生的下一秒就做出反应”。江苏正元数币的「融合智联」平台用 TDengine 做到了这件事——当室内温度越过阈值,不只是报警,而是自动触发空调切换模式。本文不按传统”选型—收益—场景”的套路展开,而是跟随一条真实的数据链路,看看 TDengine 流计算如何在数据入库的瞬间完成异常判定、告警生成,甚至反向控制设备。


1 一条数据的旅程:从设备上报到智能联动的 3 秒

先抛开架构图,我们跟随一条真实数据走一遍:

14:32:05 — 某教室温控器上报:室温 36.2℃,制冷模式,设定温度 26℃ 14:32:05 — TDengine 写入该条数据,流计算引擎在同一写入链路中判定:room_temperature > 35,触发告警规则 std000001 14:32:06 — 告警记录写入 device_alarm_record 超级表 14:32:06SceneModelTriggerProcessor 从 Redis 拉取该房间的场景策略缓存,发现策略”连续超温→强制制冷+通知管理员” 14:32:07 — 经 Redis Pipeline 批量获取关联设备信息,构建设备命令,通过 Kafka 下发至设备命令执行微服务 14:32:08 — 空调收到指令,模式切换为强制制冷,风速调至最高

从数据上报到设备响应,端到端耗时约 3 秒。这不是 demo,是正元数币「融合智联」平台的生产实践。

这 3 秒背后是一条精心设计的数据链路,我们分层拆解。


2 收益总结:TDengine 流计算引擎是核心变量

TDengine 流计算引擎带给「融合智联」的变化,不是”更快了一点”,而是从根本上改变了异常检测和联动控制的模式。

计算范式迁移:从”存后算”到”算后清”。 传统架构数据先入库、应用层定时轮询判定,数据反复搬运、延迟取决于轮询间隔。CREATE STREAM 将判定逻辑注入写入链路,数据落盘前完成计算——异常告警直接写入告警表,正常数据正常入库。一条链路同时完成存储和判定,消除了数据搬运和轮询等待。

窗口组合覆盖全场景,误报率显著降低。 COUNT_WINDOW 处理”来一条判一条”的瞬时告警,EVENT_WINDOW 处理”持续超限才告警”的持续性判定,配合 TRUE_FOR 避免抖动误报,配合 HAVING COUNT(*) >= 10 的趋势判断确保联动触发基于持续趋势而非单次读数。两种窗口组合覆盖了智慧物联几乎所有异常检测模式,应用层无需额外编写状态管理逻辑。

PRE_FILTER 预过滤 + PARTITION BY tbname 并行 + 结果物理隔离,三位一体保障性能。 正常数据在流计算入口即被过滤,不消耗计算资源;每设备独立分区并行计算,告警延迟不随设备数量增长;告警结果写入独立超级表,查询与写入互不干扰。

架构层面,流计算简化了系统、释放了研发效率。 告警规则即 SQL——新增规则只需创建 STREAM、调整阈值只需修改 PRE_FILTER 条件,不需要改 Java 代码、重新发版。应用层不再承担异常判定职责,只专注于消费告警事件、匹配策略、下发控制指令,业务逻辑更加聚焦。

业务层面,从被动监测升级为主动闭环。 设备异常从发生到告警通知 + 自动下发控制指令,端到端 3 秒完成。平台的定位从”告诉你发生了什么”升级为”帮你处理发生了什么”——异常告警 → 策略匹配 → 指令下发 → 结果回传,形成完整闭环。这种能力是传统轮询架构难以实现的。


3 数据基座:为什么不是”关系库 + 定时任务”?

在引入 TDengine 之前,「融合智联」的早期版本采用的是典型的”关系库 + 应用层轮询”方案:

  • 设备数据写入关系型数据库,每台设备一张表
  • 告警判定依赖 Java 定时任务(每 30 秒扫描一次近 1 分钟数据)
  • 联动控制依赖独立的规则引擎,通过轮询设备状态表来判断触发条件

当设备量从几百增长到几万后,问题集中爆发:

瓶颈点

具体表现

写入

万台设备每秒 1 条 = 10000 QPS 持续写入,关系库写入队列持续堆积

告警延迟

定时任务 30s 间隔 + SQL 扫描耗时 = 实际延迟常超过 2 分钟

联动响应

规则引擎轮询设备状态表,高并发下锁竞争严重

存储

3 个月历史数据超过 2TB,查询按月聚合需要数分钟

切换到 TDengine 后,核心变化发生在三个层面:

数据模型层面:一设备一子表,通过超级表统一管理,TAG 标记组织、空间、设备 ID。不再需要为每台设备手动建表,查询也无需 UNION 几百张表。

计算位置层面:TDengine 的 CREATE STREAM 将告警判定逻辑注入数据写入链路——数据落盘前先过流计算,匹配规则则直接写入告警表。应用层不再轮询,只需要订阅告警表的变化。

联动架构层面:告警产生 → 应用层消费告警事件 → 匹配场景策略 → 下发控制指令。从”定时轮询”变为”事件驱动”。


4 流计算实战:两种窗口,覆盖两类告警场景

正元数币的告警场景分为两类:瞬时越限(温度瞬间超标立即告警)和持续越限(温度持续超标一段时间才告警)。TDengine 分别用 COUNT_WINDOWEVENT_WINDOW 两种窗口类型覆盖。

图片

4.1 瞬时告警:数据入库即判定

当室内温度超过 35℃ 时,立即生成告警记录:

CREATE STREAM IF NOT EXISTS stream_alarm_dynamic_std000001 count_window(1, room_temperature)
FROM device_report_info_record
PARTITION BY tbname, organization_id, room_id, client_id
STREAM_OPTIONS(IGNORE_DISORDER | PRE_FILTER(room_temperature > 35) | MAX_DELAY(10s))
INTO device_alarm_record(ts, generate_time, alarm_rule_id, alarm_value) TAGS (
organization_id BIGINT AS organization_id,
room_id BIGINT AS room_id,
client_id VARCHAR(255) AS client_id )
AS
SELECT
CAST(_tlocaltime/1000000 AS TIMESTAMP) as ts,
_rowts as generate_time,
CAST("std000001" AS VARCHAR(10)) AS alarm_rule_id,
CAST(room_temperature AS VARCHAR(255)) AS alarm_value FROM %%tbname;

关键设计点:

  • count_window(1, room_temperature):每来 1 条数据(只要 room_temperature 非 NULL)就触发一次判定,不等攒批
  • PRE_FILTER(room_temperature > 35):在流计算入口就过滤掉正常数据,只有超限数据才进入后续逻辑,避免无效计算
  • MAX_DELAY(10s):容忍最多 10 秒的数据乱序,平衡实时性与容错

4.2 持续告警:避免瞬时抖动误报

瞬时告警的问题在于:温度可能在 35℃ 上下抖动,导致告警风暴。对于需要持续关注的场景,使用事件窗口——只有温度持续超限 30 分钟才触发告警,并输出窗口内的平均温度:

CREATE STREAM IF NOT EXISTS stream_alarm_dynamic_std000002 EVENT_WINDOW(
START WITH room_temperature > 35
END WITH room_temperature <= 35 ) TRUE_FOR (30m)
FROM device_report_info_record
PARTITION BY tbname, organization_id, room_id, client_id
STREAM_OPTIONS(IGNORE_DISORDER | MAX_DELAY(30m))
INTO device_alarm_record(ts, generate_time, alarm_rule_id, alarm_value) TAGS (organization_id BIGINT AS organization_id, room_id BIGINT AS
room_id, client_id VARCHAR(255) AS client_id) AS
SELECT
CAST(_tlocaltime/1000000 AS TIMESTAMP) as ts,
_twend as generate_time,
CAST("std000002" AS VARCHAR(10)) AS alarm_rule_id,
CAST(AVG(room_temperature) AS VARCHAR(255)) AS alarm_value FROM %%tbname
WHERE ts >= _twstart AND ts <= _twend;

EVENT_WINDOW 的精妙之处:它不是简单的”连续 N 条超限就告警”,而是追踪一个事件的完整生命周期——从温度首次突破阈值(START WITH)到回落(END WITH),如果事件持续超过 TRUE_FOR(30m),则触发告警并输出窗口内的平均值。这样既避免了抖动误报,也提供了告警期间的统计摘要。


5 用电异常实时预警:每小时用电量超标自动告警

该场景针对受控设备(如智能断路器)的每小时用电量进行实时监测,当用电量超出规则阈值时,自动生成用电异常告警记录。

图片
CREATE STREAM IF NOT EXISTS stream_electricity_alarm_dynamic_stde000001 count_window(1, electricity_consumption_hour)
FROM controlled_device_electricity_consumption_hour
PARTITION BY tbname, organization_id, room_id, controlled_device_id
STREAM_OPTIONS(PRE_FILTER(electricity_consumption_hour > 5.0) | MAX_DELAY(10s))
INTO controlled_device_electricity_consumption_hour_alarm_record( ts, generate_time, electricity_alarm_rule_id,
electricity_alarm_rule_value )
TAGS (
organization_id BIGINT AS organization_id,
room_id BIGINT AS room_id,
controlled_device_id VARCHAR(16) AS controlled_device_id )
AS
SELECT
CAST(_tlocaltime/1000000 AS TIMESTAMP) as ts,
_rowts as generate_time,
CAST("stde000001" AS VARCHAR(10)) AS electricity_alarm_rule_id, CAST(electricity_consumption_hour AS VARCHAR(255)) AS
electricity_alarm_rule_value
FROM %%tbname;

用电异常监测对数据时效性要求极高,传统方式依赖定时任务批量扫描历史数据,告警常滞后数十分钟甚至小时级。TDengine 流计算引擎通过 PARTITION BY tbname 按设备子表分区并行计算,每个受控设备的用电数据写入后立即触发流计算判定,异常告警在 10 秒内即可生成。同时,流计算结果直接写入独立的告警超级表,与应用数据物理隔离,告警查询与数据写入互不干扰。

告警记录通过定时查询同步至 MySQL 业务库:

SELECT
ts,
generate_time,
electricity_alarm_rule_id,
electricity_alarm_rule_value,
organization_id,
room_id,
controlled_device_id
FROM controlled_device_electricity_consumption_hour_alarm_record
WHERE ts >= '2025-07-25 00:00:00' AND ts < '2025-07-25 00:01:00'

6 从告警到控制:事件驱动的智能联动闭环

告警只是手段,控制才是目的。正元数币的差异化能力在于:告警不只是推送到大屏,而是自动触发设备控制。这是从”监测”到”闭环”的关键一跃。

图片

6.1 联动架构

设备上报 → TDengine(写入+流计算判定)
              ↓ 告警事件
   SceneModelTriggerProcessor(Java 微服务)
       ├── Redis 缓存:获取场景策略(避免查库)
       ├── Redis Pipeline:批量获取关联设备状态
       └── Kafka:异步下发设备命令
              ↓
      设备命令执行微服务
              ↓
         设备执行 → 结果 Kafka 回传

核心处理逻辑:

public void trigger(String clientId, Long roomId) {
    // 1. 获取该房间下有效的场景模型策略缓存
    List<SceneModelStrategyCache> validStrategies =
        sceneModelCacheHelper.getValidSceneModelStrategyCache(roomId);

    // 2. 使用 Redis 管道批量获取设备产品ID、设备ID、状态
    List<Object> devicesCommandInfoList = redisTemplate.executePipelined(...);

    // 3. 按产品ID映射设备,构建设备命令
    validStrategies.forEach(strategy -> {
           List<DeviceCommand> commandList = new ArrayList<>();
            strategy.getBelongSceneModelCache()
                .getSceneModelProductCommands()
                .forEach(productCommand -> {
                    List<DeviceCommandInfo> devices =

productDeviceCommandInfoMap.get(productCommand.getProductId());
            devices.forEach(device -> {
                commandList.add(DeviceCommand.create()
                    .clientId(device.getClientId())
                    .serviceId(productCommand.getServiceId())
                    .command(productCommand.getCommand())); });
});
       // 4. 通过 Kafka 异步下发设备命令任务
        DeviceCommandTask task = DeviceCommandTask.create()
            .deviceCommandList(commandList);
        kafkaTemplate.send(KafkaTopic.SEND_DEVICE_COMMAND_TASK, task);
    });
}

这个架构有三个设计要点:

Redis 双层缓存:场景模型策略(哪个房间、什么条件、触发什么动作)是静态配置,缓存在 Redis 中,触发时直接读取,不走数据库。设备信息同样通过 Redis Pipeline 在一次网络往返中批量获取。

Kafka 削峰填谷:当大面积温控异常(如夏季制冷系统故障)时,可能瞬间产生数千条告警,对应数千条控制命令。Kafka 作为异步消息中间件,保证命令不丢失、不阻塞主流程,设备执行微服务按自身消费能力拉取。

失败精准追溯:每条设备命令的执行结果通过 Kafka 回传,系统记录每台设备的执行状态(成功/失败/超时),失败设备可精确重试,不波及已成功的设备。

6.2 触发条件判定

并非所有告警都需要联动,部分场景需基于历史数据判断趋势。例如”房间持续低温 → 自动制热”:

SELECT DISTINCT room_id
FROM device_report_info_record
WHERE client_id IN ('device_001', 'device_002') AND ts >= NOW - 30m
AND room_temperature < 16 GROUP BY room_id
HAVING COUNT(*) >= 10;

条件是”30 分钟内同一房间有至少 10 条温度低于 16℃ 的记录”——联动触发基于持续趋势而非单次读数。


7 其它应用场景

除了上述基于流计算的实时告警与联动场景,「融合智联」平台还广泛支持”绿色校园”三大核心方向——智慧温控、水电控一体化、能耗监测。以下为各方向的关键功能摘要:

教室智慧温控、照明与业务联动

  • 场景:根据课表、作息、节假日等时间维度自动调控教室空调与照明设备,避免”长明灯””空调空转”,让师生拥有更舒适的教学环境
  • 功能:定时策略、场景模型建设(人体感应调控设备)、课表联动

宿舍安全用能管控

  • 场景:宿舍用电精细化计量与收费功能,同时集成大功率电器智能识别机制,当检测到违规用电时,可及时启动电路保护,确保用电安全。
  • 功能:用电分析、电控策略(最大功率、分时限流、过温跳电、欠费断电、恶性负载)、电费缴纳、电费缴纳(电费定价、充值缴费、免费额度设置、余额提醒)

能耗监测

  • 场景:对的空调、照明等设备进行用电分类分项计量与监测,实时掌握能耗构成与变化趋势,及时发现能耗异常点,为节能降耗决策提供数据支撑,避免因设备老化、管理疏漏或线路故障导致的能源浪费。
  • 功能:用电分析、精准计量、用电趋势对比(同比/环比/同区域对标)、能耗报告、电控策略、异常告警

关于正元数币

江苏正元数币智慧科技有限公司依托正元智慧集团(股票代码:300645)在智慧化服务领域近30年的技术积累,构建了”数字基座+行业应用服务”的产品框架体系「融合智联」,整合物联网、区块链、 5G通信、人工智能等新兴技术,提供覆盖智能终端、数据中台、聚合支付等系统的综合解决方案。其业务范围涵盖智慧校园一体化信息管理、后勤投资运营及数字人民币场景建设,并延伸至政府、军警、企业等智慧园区场景。在数字人民币领域,公司正通过子钱包开立和数字身份证技术,推动支付场景创新,已处于行业应用前沿。