引言
物联网数据管道的第一直觉是「把数据存起来」,但真正的难点在于三个矛盾。第一是写入速率与查询模式的不匹配:设备每秒上报几十万个点,而看板要的是「最近一小时的均值」,写入是点式追加,查询是范围聚合,两者对存储结构的要求截然相反。第二是数据的乱序与迟到:网络抖动、网关缓存补传、设备时钟漂移,都会让数据到达时早已过了它的时间窗口。第三是成本与保留期的矛盾:原始精度存一年,存储成本会失控,但降采样又会丢掉尖峰。
一个典型的中等规模场景是 1 万台设备、每台每分钟上报 1 条、每条 10 个字段,那就是每秒约 1.7 万个点、每天 14 亿个点。不做压缩和分层保留,这个量级很快会把单机磁盘撑爆。而时序数据库之所以存在,就是因为它用列式存储、时间戳差分编码、标签索引这套组合,把同样的数据压到关系库的十分之一甚至更小。
本文按「管道全景 → 采集 → 写入 → 模型 → 选型 → 压缩 → 降采样 → 查询 → 容量 → 告警」的顺序展开,最后给出高可用、迁移演进与常见坑。数据在边缘侧的预处理可结合 边缘计算与网关 一起看,平台侧的数据消费可结合 数字孪生与物联网平台 。
目录
- 端到端管道全景
- 采集侧接入方式
- 写入路径的工程问题
- 时序数据模型设计
- 存储选型对比
- 压缩与编码
- 降采样与保留策略
- 查询模式与 SQL 示例
- 容量估算公式
- 告警与可视化
- 高可用与容灾
- 迁移与演进
- 权衡取舍
- 常见坑清单
- 小结
1. 端到端管道全景
一条完整的物联网数据管道分为六段,每段职责清晰、可独立替换:
- 设备侧:传感器采集、边缘预处理、协议封装,产出 MQTT/CoAP 报文。
- 接入层:Broker 负责连接管理与鉴权,把消息扇出给后端。
- 处理层:规则引擎或流处理做过滤、清洗、富化、窗口聚合。
- 缓冲层:Kafka 或 Pulsar 削峰填谷,解耦写入速率与消费速率。
- 存储层:时序数据库持久化原始点与聚合结果。
- 应用层:查询 API、告警引擎、Grafana 看板、开放数据服务。
分段的目的是解耦:Broker 挂了不影响已缓冲的数据,存储慢了不会反压设备。整体架构可参考 物联网架构总览 。每段之间用明确的 schema 契约连接,避免上游改字段下游全线崩。
每一段都要有独立的可观测指标:接入层看连接数与消息速率,处理层看消费延迟与错误率,存储层看写入延迟与磁盘水位,应用层看查询 P99。哪一段先到瓶颈,扩容就加在哪一段,而不是盲目加机器。
2. 采集侧接入方式
从设备数据到处理层,有几种主流接入方式,选型取决于设备协议与运维能力。
| 方式 | 适用场景 | 优点 | 缺点 |
|---|---|---|---|
| MQTT 订阅 | 设备直连 Broker | 通用、支持 QoS | 需自己写消费者 |
| Kafka Connect | 已有 Kafka 生态 | 免代码、可扩展 | 需维护连接器 |
| Telegraf | 采集系统与网关指标 | 插件丰富 | 复杂转换能力弱 |
| Prometheus remote write | 指标类监控 | 生态成熟 | 只适合指标模型 |
| EMQX 规则引擎 bridge | MQTT 直连数据库 | 零代码、低延迟 | 复杂逻辑受限 |
| 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 | 中 | 高 | PromQL | recording 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 环形缓冲),恢复后补传。
不同场景的容灾目标差异很大:
| 场景 | RPO | RTO | 方案 |
|---|---|---|---|
| 消费级 | 分钟级 | 小时级 | 单机 + 每日备份 |
| 工业监测 | 秒级 | 分钟级 | 多副本 + Kafka 重放 |
| 计费/安防 | 接近 0 | 秒级 | 同步副本 + 双活 |
要明确 RPO 与 RTO:工业场景可能要求 RPO 接近 0,消费级可以容忍几分钟。别为了「永不丢数据」把成本推到不可接受。
12. 迁移与演进
数据管道的演进通常是「先能用,再好用」。务实路径:
- 单机 InfluxDB + 直连写入,快速验证。
- 引入 Kafka 解耦,加流处理做清洗。
- TSDB 集群化,加分层保留与连续聚合。
- 冷数据归档对象存储,热数据上 SSD。
- 引入湖仓做跨主题分析,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 加合理的保留策略就够用,过早集群化反而是负担。
下一步建议从两条线深入:边缘侧看 边缘计算与网关 把预处理与本地缓存做扎实,平台侧看 数字孪生与物联网平台 把数据消费与业务建模打通。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。