时序数据管道与存储选型

本文系统讲解物联网时序数据管道与存储选型,从设备到看板的全链路出发,覆盖 MQTT 与 Kafka Connect 采集、批量攒批与乱序迟到处理、watermark 与幂等去重、时序数据模型与标签基数控制、InfluxDB 与 TimescaleDB 与 ClickHouse 选型对比、Gorilla 与 ZSTD 压缩、连续聚合与分层保留、容量估算与告警可视化,并给出 SQL 示例与常见坑清单。

引言

物联网数据管道的第一直觉是「把数据存起来」,但真正的难点在于三个矛盾。第一是写入速率与查询模式的不匹配:设备每秒上报几十万个点,而看板要的是「最近一小时的均值」,写入是点式追加,查询是范围聚合,两者对存储结构的要求截然相反。第二是数据的乱序与迟到:网络抖动、网关缓存补传、设备时钟漂移,都会让数据到达时早已过了它的时间窗口。第三是成本与保留期的矛盾:原始精度存一年,存储成本会失控,但降采样又会丢掉尖峰。

一个典型的中等规模场景是 1 万台设备、每台每分钟上报 1 条、每条 10 个字段,那就是每秒约 1.7 万个点、每天 14 亿个点。不做压缩和分层保留,这个量级很快会把单机磁盘撑爆。而时序数据库之所以存在,就是因为它用列式存储、时间戳差分编码、标签索引这套组合,把同样的数据压到关系库的十分之一甚至更小。

本文按「管道全景 → 采集 → 写入 → 模型 → 选型 → 压缩 → 降采样 → 查询 → 容量 → 告警」的顺序展开,最后给出高可用、迁移演进与常见坑。数据在边缘侧的预处理可结合 边缘计算与网关 一起看,平台侧的数据消费可结合 数字孪生与物联网平台 。

目录

  1. 端到端管道全景
  2. 采集侧接入方式
  3. 写入路径的工程问题
  4. 时序数据模型设计
  5. 存储选型对比
  6. 压缩与编码
  7. 降采样与保留策略
  8. 查询模式与 SQL 示例
  9. 容量估算公式
  10. 告警与可视化
  11. 高可用与容灾
  12. 迁移与演进
  13. 权衡取舍
  14. 常见坑清单
  15. 小结

1. 端到端管道全景

一条完整的物联网数据管道分为六段,每段职责清晰、可独立替换:

  • 设备侧:传感器采集、边缘预处理、协议封装,产出 MQTT/CoAP 报文。
  • 接入层:Broker 负责连接管理与鉴权,把消息扇出给后端。
  • 处理层:规则引擎或流处理做过滤、清洗、富化、窗口聚合。
  • 缓冲层:Kafka 或 Pulsar 削峰填谷,解耦写入速率与消费速率。
  • 存储层:时序数据库持久化原始点与聚合结果。
  • 应用层:查询 API、告警引擎、Grafana 看板、开放数据服务。

分段的目的是解耦:Broker 挂了不影响已缓冲的数据,存储慢了不会反压设备。整体架构可参考 物联网架构总览 。每段之间用明确的 schema 契约连接,避免上游改字段下游全线崩。

每一段都要有独立的可观测指标:接入层看连接数与消息速率,处理层看消费延迟与错误率,存储层看写入延迟与磁盘水位,应用层看查询 P99。哪一段先到瓶颈,扩容就加在哪一段,而不是盲目加机器。

2. 采集侧接入方式

从设备数据到处理层,有几种主流接入方式,选型取决于设备协议与运维能力。

方式适用场景优点缺点
MQTT 订阅设备直连 Broker通用、支持 QoS需自己写消费者
Kafka Connect已有 Kafka 生态免代码、可扩展需维护连接器
Telegraf采集系统与网关指标插件丰富复杂转换能力弱
Prometheus remote write指标类监控生态成熟只适合指标模型
EMQX 规则引擎 bridgeMQTT 直连数据库零代码、低延迟复杂逻辑受限
OPC UA 采集工业设备语义完整需网关转换

实践中常见组合是:设备走 MQTT 到 EMQX,用规则引擎把遥测直接 bridge 到 Kafka,再由流处理写入 TSDB。这样设备侧无感知,处理逻辑集中在平台侧。

接入方式的选择还要考虑消息可靠性与资源占用:

  • QoS 0:最多一次,适合高频非关键指标(如心跳),丢了不补。
  • QoS 1:至少一次,适合绝大多数遥测,配合幂等去重。
  • QoS 2:恰好一次,握手开销大,仅在计费等场景使用。
  • 保留消息与遗嘱:设备离线时用 LWT 广播离线事件,让平台及时标记状态。

一个容易忽略的点是:QoS 1 会产生重复消息,如果下游不做幂等,TSDB 里就会出现重复点。所以「至少一次 + 幂等去重」是端到端可靠性的标准配置。

3. 写入路径的工程问题

写入是时序管道最容易出问题的一段,几个关键工程点:

  • 批量攒批:单点写入开销大,攒批 500 到 5000 点一批能显著提升吞吐。批太小写放大,批太大会增加延迟与内存占用。
  • 乱序与迟到:数据到达顺序不等于时间顺序。用 watermark 定义「多久之后认为窗口关闭」,迟到数据要么丢弃要么进侧输出。
  • 精确一次 vs 至少一次:精确一次需要幂等写 + 事务,成本高;至少一次配合幂等去重键 device_id + ts 更实用。
  • 幂等去重:TSDB 层用相同主键覆盖写(upsert),重复上报不会产生重复点。
  • 背压:下游慢时要能反压上游或落盘缓冲,不能让内存队列无限增长。

一个用 Python 做攒批写入的骨架,配合批大小与超时双触发:

import time

class BatchWriter:
    def __init__(self, client, batch_size=1000, flush_ms=500):
        self.client = client
        self.batch_size = batch_size
        self.flush_ms = flush_ms
        self.buf = []
        self.last = time.time()

    def add(self, point):
        self.buf.append(point)
        if len(self.buf) >= self.batch_size or (time.time() - self.last) * 1000 >= self.flush_ms:
            self.flush()

    def flush(self):
        if not self.buf:
            return
        self.client.write_points(self.buf, batch_size=len(self.buf))
        self.buf.clear()
        self.last = time.time()

三种投递语义的代价对比:

语义数据丢失重复实现成本适用
最多一次可能无最低心跳、非关键指标
至少一次无可能低多数遥测
恰好一次无无高(事务)计费、指令确认

多数场景选「至少一次 + 幂等写」,把去重交给存储层的主键覆盖,比端到端事务简单得多。

4. 时序数据模型设计

时序模型的核心是「tag 是索引,field 是数值」。以 InfluxDB 为例,measurement 相当于表,tag 是带索引的维度(device_id、region、model),field 是实际测量值(temperature、voltage)。

  • tag 与 field 的区别:tag 用于 WHERE 过滤与 GROUP BY,field 用于聚合计算。把高频变化的数值放进 tag 会制造海量序列。
  • 标签基数爆炸:如果给每条数据都带一个唯一的 request_id 作为 tag,序列数会爆炸,内存与索引全部崩掉。设备场景下,device_id 基数就是设备数,可控;但 session_id、trace_id 这类绝不能当 tag。
  • 时间精度:统一用毫秒或微秒,混用会在聚合时错位。IoT 场景毫秒通常够用,高频振动监测需要微秒。
  • 单位统一:温度统一摄氏度、功率统一瓦特,在边缘侧归一化,避免下游要按设备换算。

基数控制的一条经验法则:单个 measurement 的 tag 组合数(series cardinality)控制在百万级以内,超过就要考虑拆分 measurement 或把维度移到 field。

以 InfluxDB 的行协议(line protocol)为例,一条遥测长这样:

telemetry,device_id=dev-001,region=cn-north,model=S3 temperature=23.5,voltage=3.71 1728288000000
telemetry,device_id=dev-002,region=cn-north,model=S3 temperature=24.1,voltage=3.69 1728288000000

逗号前是 measurement,等号左侧是 tag key,右侧是 field,末尾是毫秒时间戳。注意 field 的值必须带类型(浮点、整数、字符串用引号),写错类型会导致写入失败。

5. 存储选型对比

没有银弹,选型取决于写入量、查询模式、运维能力与许可约束。

数据库写入吞吐压缩比查询语言降采样运维复杂度许可
InfluxDB 1.8高高InfluxQL连续查询低MIT
InfluxDB 3.x很高很高SQL/Flux原生中部分商用
TimescaleDB高中高完整 SQL连续聚合中Apache2/TSL
TDengine 3.x极高极高SQL原生窗口低AGPL
ClickHouse极高高SQL物化视图中高Apache2
IoTDB高极高SQL原生中Apache2
VictoriaMetrics极高极高MetricsQL下采样低Apache2
Prometheus+Thanos中高PromQLrecording rule高Apache2

选择建议:中小规模且要快速上线,InfluxDB 或 TDengine;要复杂 SQL 与关系表 join,TimescaleDB 或 ClickHouse;纯指标监控,VictoriaMetrics;工业设备树模型,IoTDB。许可上要特别注意 TDengine 的 AGPL 与 InfluxDB 3.x 的商用条款。

选型时可以按这几个问题快速收敛:数据量是否超过单机(超过 50 万点/秒考虑集群)、是否需要与业务库 join(需要则排除纯 TSDB)、团队是否熟悉 SQL(不熟则选 PromQL/Flux 反而更累)、是否能接受 AGPL(不能则排除 TDengine 社区版)。多数团队最终落在 TimescaleDB 或 ClickHouse 上,前者运维简单,后者扩展性与分析能力更强。

6. 压缩与编码

时序数据的压缩靠三件套:时间戳差分、数值 XOR、列存 + 通用压缩。

  • Delta-of-delta:时间戳通常是等间隔的,先做一阶差分得到近似常量,再对差分做二阶差分,大量为 0 或小值,编码后每点只占 1 到 2 bit。
  • XOR(Gorilla):浮点值相邻变化小,按位异或后前导零与尾随零很多,用变长编码可把 64 bit 压到十几 bit。
  • 列存 + ZSTD:把同一列连续存放再套 ZSTD/LZ4,利用列内相关性,通用压缩率再上一层。

实测压缩比:Gorilla 类编码在稳定信号上能到 10:1 到 30:1,ZSTD 在列存上再加 2 到 3 倍。也就是说 1 万个点如果原始 8 字节一个,原始 80 KB,压缩后可能只有 3 到 8 KB。压缩比高度依赖数据平稳度:剧烈波动的信号压缩比会掉到 3:1 甚至更低。

7. 降采样与保留策略

保留策略的核心是「分辨率随时间衰减」,即分层存储。

  • 热层:最近 30 天,原始精度,SSD,支撑实时查询与排障。
  • 温层:1 到 12 个月,1 分钟聚合,HDD,支撑趋势分析。
  • 冷层:1 年以上,1 小时或 1 天聚合,对象存储,支撑报表与合规。

降采样通过连续聚合(TimescaleDB 的 continuous aggregate)或物化视图实现:原始数据写入后,后台任务按窗口计算均值、最大值、最小值、计数。保留原始 + 聚合两条线,聚合线上还能做「二次聚合」(小时级由分钟级再聚合),避免重复扫描原始数据。

TTL 删除要按分区或 chunk 粒度做,避免逐行删。InfluxDB 的 retention policy、TimescaleDB 的 drop_chunks、ClickHouse 的 TTL 都是分区级删除,代价低。

三层保留策略的一个参考配置:

层级时间范围精度介质用途
热0 到 30 天原始(秒级)SSD实时看板、排障
温30 天到 1 年1 分钟聚合HDD趋势分析、周报
冷1 年以上1 小时聚合对象存储报表、合规归档

注意聚合不是简单采样,而是保留 avg、max、min、count 四组统计量:只看均值会丢掉尖峰,只看最大值又会被瞬时毛刺误导,四者一起才能在降采样后仍然可解释。

8. 查询模式与 SQL 示例

物联网查询有四类高频模式,各自的优化点不同。

第一类是时间窗口聚合:按 5 分钟求平均温度。

SELECT time_bucket('5 minutes', ts) AS bucket,
       device_id,
       avg(temperature) AS avg_temp,
       max(temperature) AS max_temp
FROM telemetry
WHERE ts >= now() - interval '1 day'
GROUP BY bucket, device_id
ORDER BY bucket DESC;

第二类是 gap filling:设备离线时段没有数据,用 time_bucket_gapfill 补 null 或插值。

SELECT time_bucket_gapfill('5 minutes', ts) AS bucket,
       locf(avg(temperature)) AS temp
FROM telemetry
WHERE device_id = 'dev-001'
  AND ts >= now() - interval '6 hours'
GROUP BY bucket;

第三类是最新值查询(last point),用 ORDER BY ts DESC LIMIT 1 或专用的 last 缓存,避免全表扫。

SELECT DISTINCT ON (device_id) device_id, ts, temperature
FROM telemetry
WHERE ts >= now() - interval '1 hour'
ORDER BY device_id, ts DESC;

第四类是设备维度下钻:按 region 或 model 聚合。这类查询靠 tag 索引加速,也最容易被高基数 tag 拖垮。它的数据模型与可观测性里的日志、指标、追踪同源,都是「时间戳 + 标签 + 数值」的骨架。

长周期查询要打到聚合表而不是原始表,这是性能的分水岭:

SELECT bucket, device_id, avg_temp
FROM telemetry_1h        -- 1 小时物化视图,行数只有原始的 1/3600
WHERE bucket >= now() - interval '90 days'
  AND region = 'cn-north'
ORDER BY bucket DESC
LIMIT 1000;

9. 容量估算公式

容量估算先算点数,再算字节,最后除以压缩比。

公式:日存储量 = 点数/秒 × 86400 × 每点字节 ÷ 压缩比。

以 1 万台设备、每分钟 1 条、每条 10 个 field 为例:

  • 点数/秒 = 10000 × 1 ÷ 60 ≈ 167 点/秒。
  • 未压缩每点约 10 field × 8 字节 = 80 字节,加时间戳与标签约 100 字节。
  • 日原始量 = 167 × 86400 × 100 ≈ 1.44 GB/天。
  • 按 10:1 压缩,实际约 144 MB/天,一年约 52 GB。

如果设备改成每秒 1 条,点数放大 60 倍,一年原始就是 3 TB 级,压缩后 300 GB,这时候就必须上分层保留。估算时还要预留索引、副本与写入放大的开销,通常在原值上乘 1.5 到 2 的安全系数。

换一个更大规模验证这个公式:10 万台设备、每秒 1 条、每条 20 个 field。点数/秒 = 10 万,每点约 200 字节,日原始量 = 100000 × 86400 × 200 ≈ 1.7 TB/天。按 10:1 压缩后 170 GB/天,一年 62 TB。这个量级单机无论如何放不下,必须集群 + 分层 + 冷归档,而且要在设计阶段就规划,不能等磁盘满了再补。

10. 告警与可视化

告警链路是「规则 → 评估 → 通知」。规则分两类:阈值告警(温度 > 80 持续 5 分钟)和异常检测(偏离基线、突增突降)。评估引擎要支持滑动窗口与去抖,否则一次抖动就发一条通知。

  • Grafana 数据源:接 InfluxDB、TimescaleDB、ClickHouse 都很成熟,用变量做设备下拉。
  • 告警规则:Grafana Alerting 或独立规则引擎,规则与看板分离,避免看板改了告警跟着变。
  • 告警疲劳:所有告警都发到同一个群,很快没人看。要分级(P0 电话、P1 群、P2 日报)、按设备聚合、支持静默与抑制。

一条 Grafana 告警规则的关键字段:for 控制去抖时长,labels.severity 决定路由级别,annotations 承载通知内容。

apiVersion: 1
groups:
  - orgId: 1
    name: iot-telemetry
    rules:
      - uid: temp-high
        title: 温度过高
        condition: C
        for: 5m                        # 持续 5 分钟才触发,抑制抖动
        data:
          - refId: A
            datasourceUid: timescale
            model:
              rawSql: "SELECT avg(temperature) AS v FROM telemetry WHERE ts > now() - interval '5 min'"
        labels:
          severity: P1
        annotations:
          summary: "设备温度超过阈值,请检查散热"

可视化设计上,趋势用折线、分布用热力图、设备状态用状态面板。Grafana 的深入用法参见 Grafana 可视化 。数据仓库侧若要做跨主题分析,可参考湖仓一体的分层建模思路,把时序明细与业务维度解耦。

11. 高可用与容灾

时序库的高可用有几条务实路线:

  • 副本:集群版 TSDB 提供多副本,写多份读一份,容忍单节点故障。
  • 异地容灾:跨可用区部署,或把聚合结果定期导出到对象存储做冷备。
  • 缓冲兜底:Kafka 保留 3 到 7 天,TSDB 故障恢复后从 Kafka 重放,数据不丢。
  • 降级:存储不可用时,边缘网关本地缓存(如 SQLite 环形缓冲),恢复后补传。

不同场景的容灾目标差异很大:

场景RPORTO方案
消费级分钟级小时级单机 + 每日备份
工业监测秒级分钟级多副本 + Kafka 重放
计费/安防接近 0秒级同步副本 + 双活

要明确 RPO 与 RTO:工业场景可能要求 RPO 接近 0,消费级可以容忍几分钟。别为了「永不丢数据」把成本推到不可接受。

12. 迁移与演进

数据管道的演进通常是「先能用,再好用」。务实路径:

  1. 单机 InfluxDB + 直连写入,快速验证。
  2. 引入 Kafka 解耦,加流处理做清洗。
  3. TSDB 集群化,加分层保留与连续聚合。
  4. 冷数据归档对象存储,热数据上 SSD。
  5. 引入湖仓做跨主题分析,TSDB 只保留近期数据。

迁移时最大的坑是 schema 不兼容:老数据的 tag 与 field 定义和新数据不一致。解决办法是在写入层做 schema 归一化,并为历史数据做一次性回填。双写过渡期要保证幂等,避免重复点。

双写过渡的具体做法:新老库同时写入,查询侧先切读新库做灰度,观察一段时间(至少覆盖一个完整的业务周期,如一周),确认数据一致后再停掉老库写入。回滚方案要提前准备好,否则灰度期间出问题会手忙脚乱。

权衡取舍

决策点选项 A选项 B建议
写入语义精确一次至少一次 + 幂等后者性价比更高
原始数据保留全量长期降采样 + 分层分层,热 30 天
标签设计维度全进 tag高基数进 field高基数别当 tag
存储单机 TSDB集群 + 湖仓按规模,别过度设计
压缩追求极致比兼顾查询速度查询延迟优先

常见坑清单

  • 现象:TSDB 内存暴涨后 OOM。原因:高基数 tag 制造海量序列。规避:审计 tag 基数,高基数维度移入 field。
  • 现象:看板查询超时。原因:全表扫原始数据算长周期聚合。规避:用连续聚合/物化视图预计算。
  • 现象:数据重复计数。原因:至少一次投递 + 无幂等键。规避:以 device_id + ts 为幂等主键覆盖写。
  • 现象:凌晨数据缺失。原因:设备时钟漂移导致迟到数据被丢。规避:watermark 放宽 + 迟到侧输出。
  • 现象:磁盘一周写满。原因:无保留策略,原始数据无限增长。规避:retention policy + 分区级 TTL 删除。
  • 现象:写入吞吐上不去。原因:单点写入、无攒批。规避:批大小 500 到 5000,双触发 flush。
  • 现象:Gap 处曲线断裂误判。原因:设备离线无数据,未做 gap filling。规避:time_bucket_gapfill + locf。
  • 现象:告警风暴。原因:无去抖、无分级。规避:滑动窗口 + 聚合 + 分级通知。
  • 现象:迁移后历史数据查不到。原因:新旧 schema 不一致。规避:写入层归一化 + 一次性回填。

小结

时序数据管道的本质是把「高频写入、低频聚合查询」这对矛盾用分层设计化解:Broker 解耦连接、Kafka 削峰、流处理清洗、TSDB 存储、连续聚合加速查询、分层保留控制成本。每一段都能独立演进,不必一次做到位。

选型上不要陷入「哪个数据库最好」的争论,先算清楚点数、字节、保留期和查询模式,再对照对比表选。多数中等规模场景,单机 TSDB 加合理的保留策略就够用,过早集群化反而是负担。

下一步建议从两条线深入:边缘侧看 边缘计算与网关 把预处理与本地缓存做扎实,平台侧看 数字孪生与物联网平台 把数据消费与业务建模打通。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「物联网」更多文章

  1. 工业物联网协议与网关
  2. 边缘 AI 推理
  3. 设备配网与批量运维