TDengine 在设备监测中的实际落地:从传感器上报到实时看板

设备监测场景里,烟感、温湿度、气体、水浸、门磁等数据会持续上报。业务通常还要支持:

  • 查看某台设备的最新状态;
  • 按位置展示实时环境数据;
  • 查询一段时间内的变化曲线;
  • 按区域和告警类型做统计;
  • 设备异常时触发通知和抓拍。

这些数据如果全部写进 MySQL,前期实现很快,数据量上来后,明细表、索引、统计查询和备份都会越来越重。

本文结合一套真实的设备监测链路,重点讲清楚:

  1. 不同传感器数据怎么统一写入 TDengine;
  2. 超级表、子表、Tag 和 Field 怎么拆;
  3. Schemaless 写入怎么封装;
  4. 最新数据和告警统计怎么查询;
  5. 实际代码里有哪些值得保留和改进的点。

文中名称均已通用化,Java 时间戳统一使用毫秒


一、先看完整链路

项目中的时序数据不是单独存在的,它从设备接入开始,最终服务于实时看板和告警统计。

设备上报

协议解析

标准化传感器记录

按设备类型转换

烟感记录

气体记录

温湿度记录

水浸 / 光栅 / 门磁记录

TsDB 写入工具

TsModel 拆分 Tag / Field

LINE 协议

TDengine

最新数据查询

时间窗口查询

告警聚合

实时看板

统计大屏

实际写入可以概括为:

设备原始消息
  → 解析成统一传感器记录
  → 按设备类型转换成对应的 TsModel
  → TsDBUtils 提取表名、Tag、Field
  → TimeSerial 组装 LINE
  → SchemalessWriter 写入 TDengine

查询则是反方向:

Mapper 查询 TDengine
  → 获取最新记录 / 时间窗口数据 / 告警聚合
  → 业务库补充设备名称、位置、抓拍图片
  → 组装实时看板或告警列表

二、为什么不同传感器要拆成不同超级表

烟感、温湿度和气体看起来都属于“环境数据”,但它们的数据结构并不相同。

类型主要变化数据
烟感是否报警
水浸是否报警
门磁开关或报警状态
温湿度温度、湿度、上下限、判定结果
气体多种气体指标或原始 JSON

如果全部塞进一张万能表,会出现大量空字段,字段类型和查询逻辑也会越来越混乱。

更合适的设计是:一个清晰的数据领域对应一个超级表。

smoke_record                   烟感明细
temperature_humidity_record    温湿度明细
gas_record                     气体明细
immersion_record               水浸明细
door_magnet_record             门磁明细
device_alarm_record            统一告警事件

这样做有三个好处:

  1. 同类数据结构稳定;
  2. 查询不需要判断大量空字段;
  3. 新增设备类型时,不会污染已有超级表。

结构化列和 JSON 怎么选

实际场景里还会遇到另一种取舍:

  • 温度、湿度这类稳定指标,拆成数值列;
  • 气体设备的测点可能因型号不同而变化,可以先保留原始 JSON。

稳定指标拆列后,可以直接做比较、聚合和降采样:

SELECT AVG(temperature), MAX(humidity)
FROM temperature_humidity_record
WHERE ts >= ? AND ts < ?;

JSON 的优点是兼容型号差异,缺点是统计前需要解析,压缩、类型校验和查询效率通常也不如数值列。

一个更实用的折中方案是:

常用指标 → 独立数值 Field
完整报文 → raw_payload 字符串 Field

这样既能直接查询核心指标,也保留了原始数据,便于设备协议升级后追溯。


三、写入第一步:把业务对象转换成时序模型

设备解析层先得到一份统一记录,再根据设备类型转换成对应的时序模型。

以气体数据为例,实际思路很简单:

public void saveGasDataToTs(SensorRecord source) {
    GasRecordTs target = new GasRecordTs();

    // 复用通用设备、位置和测点信息
    BeanUtils.copyProperties(source, target);

    // 当前实现使用平台接收时间作为时序时间
    long timestampMillis = System.currentTimeMillis();

    TimeSeriesUtils.save(target, timestampMillis);
}

烟感、水浸、门磁等开关量设备会额外从原始 JSON 中提取报警状态:

public void saveSmokeDataToTs(SensorRecord source) {
    SmokeRecordTs target = new SmokeRecordTs();
    BeanUtils.copyProperties(source, target);

    int alarm = parseInt(
            source.getSensorData().path("alarm").asText(),
            AlarmState.NORMAL);

    target.setAlarm(alarm);
    TimeSeriesUtils.save(target, System.currentTimeMillis());
}

采集时间还是接收时间

当前这类写法使用平台接收时间,优点是:

  • 不受设备时钟错误影响;
  • 所有数据使用统一时钟;
  • 适合实时状态和告警展示。

但如果设备支持离线缓存和补传,只使用接收时间会把历史数据挤在补传时刻。

更完整的模型可以同时保存两个时间:

public record SensorEvent(
        long collectedAtMillis, // 设备采集时间
        long receivedAtMillis,  // 平台接收时间
        JsonNode payload) {
}

通常用 collectedAtMillis 作为 TDengine 主时间戳,receivedAtMillis 作为 Field,用于排查网络延迟和设备时钟漂移。设备不可信时,再退化为平台接收时间。


四、用 TsModel 统一约束 Tag 和 Field

不同传感器模型虽然字段不同,但都可以实现同一个契约:

public interface TsModel {

    /** 子表级、相对稳定的维度 */
    Map<String, String> toTagMap();

    /** 每条记录真正变化的数据 */
    Map<String, String> toFieldMap();
}

以烟感记录为例:

@TableName("smoke_record")
public class SmokeRecordTs implements TsModel {

    private long id;
    private long deviceId;
    private long locationId;
    private int deviceType;
    private int alarm;

    @Override
    public Map<String, String> toTagMap() {
        return Map.of(
                "device_id", String.valueOf(deviceId),
                "location_id", String.valueOf(locationId),
                "device_type", String.valueOf(deviceType)
        );
    }

    @Override
    public Map<String, String> toFieldMap() {
        return Map.of(
                "id", DataType.BIGINT(id),
                "alarm", DataType.INT(alarm)
        );
    }
}

温湿度模型则把温度、湿度和判定结果放进 Field:

@Override
public Map<String, String> toFieldMap() {
    return Map.of(
            "id", DataType.BIGINT(id),
            "temperature", DataType.FLOAT(temperature),
            "humidity", DataType.FLOAT(humidity),
            "temperature_result", DataType.INT(temperatureResult),
            "humidity_result", DataType.INT(humidityResult)
    );
}

Tag 和 Field 怎么选

可以按下面的标准判断:

问题
创建子表后基本不变化?可以考虑 Tag更适合 Field
经常用于过滤或分组?可以考虑 Tag更适合 Field
每条记录都可能不同?不应作为 Tag放 Field

设备 ID 虽然基数很高,但如果目标就是“一台设备一张子表”,它仍然适合作为 Tag。

设备名称、负责人、实时在线状态等容易变化,不适合放 Tag。查询时根据设备 ID 从业务库补充更稳妥。


五、通用写入工具做了什么

业务 Service 不应该关心 LINE 协议细节。通用工具可以统一完成:

  1. @TableName 获取超级表名;
  2. TsModel 获取 Tag 和 Field;
  3. 检查 Tag 与 Field 是否重名;
  4. 组装 Schemaless LINE;
  5. 按毫秒精度写入 TDengine。
public final class TimeSeriesUtils {

    private TimeSeriesUtils() {
    }

    public static <T extends TsModel> void save(T target,
                                                long timestampMillis) {
        TableName table = target.getClass().getAnnotation(TableName.class);
        if (table == null) {
            throw new IllegalArgumentException("时序模型缺少 @TableName");
        }

        Map<String, String> tags = target.toTagMap();
        Map<String, String> fields = target.toFieldMap();

        validateNoDuplicateKeys(tags, fields, table.value());

        TimeSerial.getInstance("device_data")
                .add(
                        table.value(),
                        tags,
                        fields,
                        timestampMillis)
                .write();
    }

    private static void validateNoDuplicateKeys(
            Map<String, String> tags,
            Map<String, String> fields,
            String table) {
        Set<String> duplicated = new HashSet<>(tags.keySet());
        duplicated.retainAll(fields.keySet());

        if (!duplicated.isEmpty()) {
            throw new IllegalArgumentException(
                    "Tag 与 Field 重名, table=" + table
                            + ", keys=" + duplicated);
        }
    }
}

Tag 和 Field 重名不是普通数据问题,而是确定性的建模错误。写入前直接失败,比等 TDengine 自动建出错误 Schema 后再迁移更划算。


六、Schemaless LINE 是怎么生成的

最终写入 TDengine 的数据类似:

smoke_record,_sub_table=smoke_record_abcd1234,device_id=1001,location_id=20 id=90001i64,alarm=1i32 1710000000000

可以拆成四部分:

超级表和子表             Tag                    Field                  毫秒时间戳
smoke_record,...         device_id=...          id=...,alarm=...       1710000000000

1. 子表名要可重复生成

当前思路是根据 Tag 组合生成摘要,让同一组 Tag 始终写入同一张子表:

private static String buildSubTable(String stable,
                                    Map<String, String> tags) {
    String canonicalTags = tags.entrySet().stream()
            .sorted(Map.Entry.comparingByKey())
            .map(entry -> entry.getKey() + "=" + entry.getValue())
            .collect(Collectors.joining("&"));

    return stable + "_" + sha256(canonicalTags);
}

Tag 必须先排序。直接对 HashMap 的遍历结果做摘要,可能因为遍历顺序变化,让相同 Tag 生成不同子表。

2. Field 类型要显式编码

Schemaless 首次写入会影响字段类型。建议统一使用类型工具:

DataType.BIGINT(1001L);  // 1001i64
DataType.INT(1);         // 1i32
DataType.FLOAT(26.5f);   // 26.5f32
DataType.DOUBLE(26.5d);  // 26.5f64
DataType.NCHAR("文本");  // NCHAR 字符串

不要让同一个 Field 有时写整数、有时写字符串。Schema 一旦被污染,后续通常只能通过新表和迁移修复。

3. 字符串需要统一转义

Tag 中的空格、逗号和等号具有语法含义:

private static String escapeTag(String value) {
    return value
            .replace("\\", "\\\\")
            .replace(" ", "\\ ")
            .replace(",", "\\,")
            .replace("=", "\\=");
}

JSON Field 也需要先处理引号和反斜杠,再交给 NCHAR 编码。业务代码不要自行拼接 LINE 字符串。


七、最新数据怎么查询

实时看板最常见的需求不是查历史曲线,而是:

给我某台设备或某个位置最近的一条有效数据。

当前常见做法是两步查询:

-- 第一步:获取最近一个 ID
SELECT MAX(id)
FROM smoke_record
WHERE device_id = ?
  AND ts >= ?;
-- 第二步:根据 ID 回表
SELECT id, device_id, location_id, alarm, ts
FROM smoke_record
WHERE id = ?;

这种写法容易理解,但有两个问题:

  1. 需要两次数据库往返;
  2. MAX(id) 代表最大 ID,不一定严格等于最新采集时间。

更直接的写法是按时间排序:

SELECT id, device_id, location_id, alarm, ts
FROM smoke_record
WHERE device_id = ?
  AND ts >= ?
ORDER BY ts DESC
LIMIT 1;

如果当前 TDengine 版本支持并且已经验证语义,也可以使用 LAST_ROW 获取完整的最后一行。不要随意把多个字段分别写成 LAST(field),因为字段存在 NULL 时,各列返回的最后非 NULL 值可能不是同一行。

为什么要加 timeFlag

“最后一条”不等于“当前有效”。

设备半年前上报过一次,今天直接把那条数据当实时状态显然不合理。因此查询最新数据时,通常需要一个有效时间下限:

ts >= 当前时间 - 有效窗口

这个窗口应按设备类型配置,而不是所有传感器共用一个固定值。


八、按位置批量查询最新数据

看板常常要同时展示多个位置。如果逐个位置查一次,会产生大量数据库往返。

更适合的思路是一次传入位置集合,在 TDengine 中分组取最后一条:

SELECT LAST_ROW(*)
FROM temperature_humidity_record
WHERE ts >= ?
  AND location_id IN (?, ?, ?)
PARTITION BY location_id;

如果版本或驱动不支持上面的写法,可以使用子查询或经过验证的 LAST() 聚合,但要确保同一位置的所有字段来自同一条记录。

业务层再根据结果组装看板:

Map<Long, TemperatureHumidityRecord> latestByLocation =
        records.stream().collect(Collectors.toMap(
                TemperatureHumidityRecord::getLocationId,
                Function.identity()));

批量查一次再映射,通常比循环调用“单位置最新一条”更稳定。

历史曲线要先降采样

最新点查询解决的是“现在怎么样”,历史曲线解决的是“一段时间怎么变化”。

如果设备每几秒上报一次,查询一个月时不应该把全部原始点直接返回前端。可以按时间窗口聚合:

SELECT
    _wstart AS bucket_start,
    AVG(temperature) AS avg_temperature,
    MAX(temperature) AS max_temperature,
    MIN(temperature) AS min_temperature
FROM temperature_humidity_record
WHERE device_id = ?
  AND ts >= ?
  AND ts < ?
INTERVAL(5m)
FILL(NULL);

窗口大小要根据查询跨度动态调整:

查询跨度建议窗口
近 1 小时1 分钟
近 1 天5~15 分钟
近 1 个月1~6 小时

FILL 也不能随便选:

  • FILL(NULL):明确表示该窗口没有数据,最安全;
  • FILL(PREV):适合短时间保持不变的状态,不适合长时间离线;
  • FILL(VALUE, 0):只适合业务上“无数据就是 0”的指标;
  • 线性插值:适合连续缓慢变化的指标,但不能用于离散告警状态。

温湿度缺数据时填 0,前端可能会误以为真的出现了 0℃。因此默认建议保留 NULL,再由展示层决定是否断线。


九、告警数据为什么要单独建表

传感器明细和告警事件的查询方式不同:

  • 明细用于曲线、最新状态和历史回放;
  • 告警用于统计数量、按类型筛选和处理闭环。

因此可以在传感器明细之外,再维护一张统一告警超级表:

device_alarm_record

代表性设计:

类型字段
Tagdevice_id、location_id、device_type、area_id
Fieldid、alarm_type、sensor_data、normal_data

这里不建议把设备名称放进 Tag:

  • 名称可能修改;
  • 名称可能包含逗号、空格等特殊字符;
  • 名称会增加子表变化风险。

告警表保存设备 ID,展示时再从业务库补充名称。

告警写入后,通知和抓拍可以继续异步处理:

传感器数据

是否命中告警规则

只保存明细

写入告警时序表

发送通知

触发摄像头抓拍


十、告警统计查询

大屏常见需求是:按区域、位置和告警类型统计数量。

SELECT
    area_id,
    alarm_type,
    COUNT(*) AS alarm_count
FROM device_alarm_record
WHERE ts >= ?
  AND ts < ?
  AND area_id IN (?, ?, ?)
GROUP BY area_id, alarm_type;

时间范围建议统一使用左闭右开:

[startTimeMillis, endTimeMillis)

也就是:

ts >= startTimeMillis AND ts < endTimeMillis

相比 BETWEEN start AND end,左闭右开更适合日报、周报和分段迁移,可以避免相邻窗口重复统计边界点。

动态查询条件可以继续使用 MyBatis:

<where>
    <if test="startTimeMillis != null and endTimeMillis != null">
        AND ts <![CDATA[>=]]> #{startTimeMillis}
        AND ts <![CDATA[<]]>  #{endTimeMillis}
    </if>
    <if test="alarmType != null">
        AND alarm_type = #{alarmType}
    </if>
</where>
GROUP BY area_id, alarm_type

十一、TDengine 只存数据,业务信息怎么补

时序表适合保存设备 ID、位置 ID 和指标,不适合冗余大量容易变化的业务信息。

告警列表通常需要:

  • 设备名称;
  • 位置名称;
  • 区域名称;
  • 抓拍图片;
  • 处理状态。

更稳妥的做法是:

先查 TDengine 告警分页
  → 收集 deviceId / locationId / alarmId
  → 批量查询业务库和图片表
  → Map 预加载
  → 组装 VO
Map<Long, DeviceInfo> deviceMap =
        deviceRepository.batchByIds(deviceIds);

Map<Long, LocationInfo> locationMap =
        locationRepository.batchByIds(locationIds);

Map<Long, CaptureInfo> captureMap =
        captureRepository.batchByAlarmIds(alarmIds);

List<AlarmVO> result = alarms.stream()
        .map(alarm -> convert(
                alarm,
                deviceMap.get(alarm.getDeviceId()),
                locationMap.get(alarm.getLocationId()),
                captureMap.get(alarm.getId())))
        .toList();

重点是批量预加载,不要在循环中逐条查询业务库,否则很容易形成 N+1。


十二、Schemaless 首次写入和“表不存在”

Schemaless 的表可能在第一条数据到来前并不存在。

查询时出现“表不存在”,它可能表示:

  1. 这个设备类型从未产生过数据;
  2. 初始化数据还没写入;
  3. 表名或迁移脚本有问题。

不能把所有异常都吞成空列表。可以区分处理:

public void handleTimeSeriesQueryException(String operation,
                                           Exception exception) {
    if (isTableNotExists(exception)) {
        log.warn("{}:时序表尚未创建", operation);
        metrics.increment("tsdb.table_not_exists", operation);
        return;
    }

    log.error("{}:时序查询失败", operation, exception);
    throw new TimeSeriesQueryException(operation, exception);
}

有些系统会在设备注册时写入一条正常态初始化数据,提前创建超级表和子表。它可以减少首次查询异常,但也要注意:

  • 初始化记录不能污染真实统计;
  • 查询曲线时要能识别或过滤占位记录;
  • 如果没有强需求,懒创建也可以,但必须做好“无表”和“无数据”的区分。

十三、写入可靠性和性能

1. 当前同步写的特点

通用工具直接调用同步 write() 时,调用线程会等待 TDengine 返回。

优点:

  • 能立即知道写入是否成功;
  • 异常处理简单。

缺点:

  • TDengine 抖动会阻塞设备消费线程;
  • 每条数据单独写,网络往返较多;
  • WebSocket 封装即使接收 List,也可能在内部逐行发送。

2. 高吞吐场景的改造方式

失败

设备消息

解析标准记录

有界缓冲队列

按表 / 子表分批

批量写 TDengine

重试 / 死信

建议:

  • 使用有界队列,避免 TDengine 变慢时耗尽 JVM;
  • 按数量或时间触发批量 flush;
  • 不要无限重试;
  • 监控队列长度、写入耗时和失败数量;
  • 普通遥测可以异步批量,关键告警继续同步或可靠消息写入。

批次可以从 100~1000 行开始压测,再根据连接方式和单行大小调整。

3. 同步写和单向写的边界

模式适用注意
同步写告警、关键统计点延迟更高,但能确认结果
单向写可重放的普通遥测不等待结果,不等于可靠

writeOneway() 只是不等待服务端响应。若要保证可靠性,仍需要上游 MQ、失败重放或本地缓冲。


十四、时间戳和重复数据

全文统一使用毫秒:

long timestampMillis = event.getCollectedAtMillis();

驱动也明确指定:

SchemalessTimestampType.MILLI_SECONDS

还需要提前确认两个问题:

1. 同一子表、同一毫秒内会不会有多条

如果设备上报频率很高,毫秒精度可能不够。需要根据实际情况:

  • 提升时间精度;
  • 合并同一毫秒的数据;
  • 或调整子表划分。

2. 设备重传怎么处理

同一子表、同一时间戳重复写入时,具体覆盖或更新行为要结合 TDengine 版本和配置验证。

业务设计里要明确:

  • 重传是否应该覆盖原值;
  • 是否保留原始报文 ID;
  • 是否需要记录接收时间;
  • 失败重试是否可能产生重复点。

十五、Schema 变更和历史迁移

Schemaless 不等于没有 Schema。Field 类型、Tag 或超级表名发生变化时,历史数据仍然需要迁移。

推荐按时间窗口迁移:

public void migrate(long startTimeMillis,
                    long endTimeMillis,
                    long windowMillis) {
    for (long start = startTimeMillis;
         start < endTimeMillis;
         start += windowMillis) {

        long end = Math.min(
                start + windowMillis,
                endTimeMillis);

        // 左闭右开,避免边界重复
        List<OldPoint> oldPoints =
                oldRepository.query(start, end);

        List<NewPoint> newPoints =
                oldPoints.stream()
                        .map(this::convert)
                        .toList();

        batchWrite(newPoints);
        recordCheckpoint(end);
    }
}

迁移任务至少要具备:

  • 可配置时间窗口;
  • 断点记录;
  • 可重复执行;
  • 新表存在性检查;
  • 迁移前后数量校验;
  • 关键字段抽样比对。

还要避免“模型写新表、Mapper 继续查旧表”这类常见遗留。表名最好由模型元数据统一提供,减少 Java 模型和 XML SQL 各写一份造成的不一致。


十六、这套设计中值得保留的点

结合实际应用,有几处设计很值得复用:

  1. 按传感器领域拆超级表,而不是做一张万能宽表;
  2. TsModel 契约化 Tag 和 Field,业务 Service 不直接拼 LINE;
  3. 写入前检查 Tag/Field 重名,让 Schema 错误尽早暴露;
  4. Field 类型统一编码,避免 Schemaless 类型漂移;
  5. 传感器明细和告警事件分开,分别服务曲线和统计;
  6. 最新数据设置有效时间窗口,避免历史旧点冒充实时状态;
  7. 告警分页后批量补充业务信息,避免 N+1;
  8. 表不存在单独降级,其他查询异常保留完整错误。

十七、仍然值得继续优化的点

  1. 平台接收时间和设备采集时间需要明确优先级;
  2. MAX(id) + 回表 可以优化为按 ts 直接取最后一行;
  3. BETWEEN>= / < 应统一为左闭右开;
  4. 设备名称等可变字符串不应作为 Tag;
  5. 同步逐条写入可以升级为有界队列和批量 flush;
  6. 子表摘要前必须固定 Tag 顺序;
  7. Java 模型与 Mapper SQL 的表名应保持单一来源;
  8. 初始化占位数据不能污染真实统计;
  9. 需要明确原始数据保留周期和迁移策略。

收个尾

TDengine 在设备监测中的落地,不只是把 MySQL 换成另一个数据库。

真正的主线是:

设备数据先标准化,再按领域转换成时序模型;TsModel 明确 Tag 和 Field;Schemaless 负责写入;Mapper 负责最新点、时间窗口和告警聚合;业务库负责补充名称、位置和处理状态。

这条链路里,建模比写入更重要,查询边界比 SQL 是否简短更重要,异常、批量和迁移则决定了它能不能长期稳定运行。

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐