文档类型:技术与架构实践分享
编写日期:2026-09-22
实践状态:SSE 实时推送与 InfluxDB 按点位间隔保存已完成现场功能验证
适用场景:设备高频遥测、实时大屏、时序数据持久化、Kafka 消费解耦
工业物联网平台经常同时面对两类目标:
本实践最终采用“同一原始遥测 Topic、两个独立 Kafka 消费组”的双链路架构:实时链路只负责最新状态投影与 SSE 推送,持久化链路独立完成按点位间隔判断、批量写入 InfluxDB、水位提交和异常重试。两个链路共享无状态的标准化规则和点位映射缓存,但不共享消费位点、消费线程或下游成败状态。
现场功能验证表明:
本次验证属于功能与链路正确性验证。文档不虚构生产吞吐量、P95/P99 延迟或容量上限,正式容量结论仍应以压测和生产监控为准。
本文所有名称和数据均为通用化示例:
<raw-topic>、<realtime-group> 等占位符;tenant-demo、device-demo、temperature 等示例值;假设设备网关每秒上报一次遥测,一条设备消息包含温度、压力、瞬时流量、累计流量、电能等多个点位。
业务侧通常有不同要求:
| 目标 | 典型要求 |
|---|---|
| 实时大屏 | 尽可能快速展示最新值,允许在一个推送窗口内合并旧值 |
| 历史趋势 | 根据点位重要性每 5 秒、30 秒或 60 秒保存一次 |
| 累计量与抄表 | 保留可靠时间戳和原始累计值,便于期间差值计算 |
| 故障恢复 | 数据库短时不可用时不能让实时页面同步停顿 |
| 性能 | 批量处理,避免每条消息、每个点位都查询关系数据库 |
这里最容易混淆的是“设备采集频率”和“平台入库间隔”:
早期链路可以抽象为:
原始遥测
→ 消息校验与点位映射
→ 入库间隔判断
→ 写 InfluxDB
→ 更新实时投影
→ SSE 推送
该结构的核心问题不是调用顺序本身,而是实时和持久化共用一个 Kafka 消费组、同一批线程和同一个 Offset 推进条件。
当 InfluxDB 变慢、写入重试、Redis 水位操作耗时或点位映射频繁回源数据库时,当前批次无法及时完成,Kafka Offset 不能推进,后续实时消息也只能排队。外部表现通常是:
一种常见方案是在原始 Topic 后增加标准化 Topic,再由实时组和持久化组分别消费。它可以建立清晰的内部消息契约,但也增加一次 Kafka 生产、存储和消费过程。若标准化生产者自身吞吐不足、需要等待 Broker 确认,或者标准化阶段频繁查询数据库,它会成为两个下游共同的前置瓶颈。
因此,本次实践没有把中转 Topic 作为当前主链路,而是保留为可回滚、可演进的兼容模式。
关键隔离点:
两条链路都会把原始设备消息转换为统一的内部对象。标准化步骤包括:
脱敏后的标准化消息示例:
{
"schemaVersion": "1.0",
"eventId": "<topic>:2:10001",
"tenantId": "tenant-demo",
"deviceId": "device-demo",
"deviceCode": "METER-DEMO-01",
"sourceTs": 1790056800000,
"ingestTs": 1790056800125,
"points": [
{
"pointCode": "active_power",
"valueType": "double",
"value": 428.6,
"unit": "kW",
"samplingIntervalSec": 5
}
]
}
eventId 可以由原 Topic、分区和 Offset 组成,便于定位和幂等处理。示例不使用真实 Topic、设备编号和 Offset。
每个设备点位的有效入库间隔按以下优先级确定:
设备点位覆盖值 samplingIntervalOverride
> 型号点位默认值 samplingIntervalSec
> null:每条合法数据都允许入库
设备采集频率字段不参与该判断。设备仍可每秒上报,SSE 仍可秒级更新,而 InfluxDB 根据点位配置每 5 秒或 60 秒保存一次。
设备覆盖表用于保存少量设备例外,例如:
| 字段 | 含义 |
|---|---|
device_id |
需要单独配置的设备 |
point_template_id |
被覆盖的型号点位 |
enabled_override |
是否启用该设备点位 |
sampling_interval_override |
设备级入库间隔,单位秒 |
unit_override |
设备级工程单位 |
effective_start/effective_end |
生效时间范围 |
该表属于遥测标准化的运行时依赖。部署前必须确保表结构、索引和应用 search_path 正确;不能把“表不存在”作为空配置静默忽略,否则两个消费组会在冷缓存加载时同时失败。
非法间隔不应导致数据被丢弃。实践中对超出允许范围的历史配置记录告警,并降级为 null,即每条合法数据均可入库。
双消费组会分别执行标准化。如果每条消息都查询设备表、型号点位表和设备覆盖表,数据库访问量会随上报频率线性增长,并成为 Kafka 消费瓶颈。
最终缓存结构为:
请求设备点位映射
→ L1:进程内 Caffeine
→ L2:Redis
→ PostgreSQL 回源
同时对“设备映射”和“型号点位集合”分别缓存:
配置变更后执行:
Redis TTL 增加随机抖动,避免大量设备同时过期形成缓存雪崩。Redis 异常时开启短时间本机旁路,直接单次回源 PostgreSQL 并写入 L1,避免每条消息反复等待同一个 Redis 故障。
需要注意:Redis 是缓存,不是设备配置的最终事实来源。PostgreSQL 仍是设备、型号和覆盖配置的权威数据源。
实时链路只做以下工作:
消费原始批次
→ 标准化
→ 按 tenant + device + point 合并为批次最新值
→ 更新实时投影
→ Redis Pub/Sub
→ SSE Hub
批次合并示例:
| 点位 | 批次内时间 | 值 | 处理结果 |
|---|---|---|---|
| pressure | 10:00:00.000 | 0.61 | 被同批次新值替代 |
| pressure | 10:00:00.500 | 0.62 | 被同批次新值替代 |
| pressure | 10:00:00.900 | 0.63 | 发布 |
这意味着 SSE 是“最新状态流”,不是完整历史审计流。它追求页面及时性和有界内存,不承诺浏览器逐条接收设备上报的每一个瞬时值。
实时链路明确不执行:
因此,把某点位入库间隔从 1 秒调整为 60 秒,不会把该点位 SSE 推送也降低为每分钟一次。
每个点位维护最后一次成功入库的源时间:
key = tenantId + deviceId + pointCode
value = lastPersistedSourceTs
判断规则:
没有成功水位 → 选择当前点
入库间隔为 null → 选择当前点
sourceTs < lastPersistedSourceTs → 判定为乱序
sourceTs >= last + interval → 选择当前点
其他情况 → 跳过本次历史入库
伪代码:
if (last != null && sourceTs < last) {
sendToPersistenceDlq(event);
} else if (last == null || interval == null
|| sourceTs >= last + interval.toMillis()) {
selected.add(event);
provisionalWatermark.put(key, sourceTs);
} else {
skipped++;
}
假设设备每秒上报,点位入库间隔为 5 秒:
| 源时间 | SSE | InfluxDB | 原因 |
|---|---|---|---|
| 12:00:00 | 更新 | 写入 | 首个样本 |
| 12:00:01 | 更新 | 跳过 | 未到5秒 |
| 12:00:02 | 更新 | 跳过 | 未到5秒 |
| 12:00:03 | 更新 | 跳过 | 未到5秒 |
| 12:00:04 | 更新 | 跳过 | 未到5秒 |
| 12:00:05 | 更新 | 写入 | 到达下一个周期 |
| 12:00:06 | 更新 | 跳过 | 未到下一周期 |
| 12:00:10 | 更新 | 写入 | 到达下一周期 |
该算法以首次成功入库时间为滚动起点,不强制对齐自然分钟或自然整点。
一个 Kafka 批次内可能包含同一点位的多条消息。判断过程中先更新“批次临时水位”,避免在真实水位尚未提交时选中多个未到期样本。
真实水位只能在 InfluxDB 写入成功后提交:
采样判断
→ 生成候选点
→ 分片批量写 InfluxDB
→ 全部写入成功
→ 更新本地水位
→ Redis Pipeline 同步水位
→ Listener 正常返回并提交 Kafka Offset
如果 InfluxDB 写入失败:
这是“宁可少量重复,不可静默少写”的可靠性取舍。
| 消费组 | 数据源 | 首次启动策略 | Offset 推进条件 |
|---|---|---|---|
| 实时组 | 原始遥测 Topic | 从最新位置开始 | 当前批次实时投影成功 |
| 持久化组 | 原始遥测 Topic | 复用既有入库组位点 | 当前批次持久化流程成功 |
实时组从最新位置开始,是因为实时页面关心“现在”,不应在首次启用时回放大量历史消息。组一旦产生提交位点,重启后继续从已提交位置恢复。
持久化组沿用旧链路的消费组 ID,可以继承原 Offset,降低切换时从头回放或跳过积压的风险。切换消费组 ID 属于数据语义变更,不能只当普通配置修改。
直接双链路模式对处理异常执行持续重试,不将临时故障伪装成成功:
持续重试保证不静默丢失,但也意味着永久性配置错误会阻塞对应分区。例如点位覆盖表缺失、字段不兼容或必需 Topic 不存在。因此必须同时建设启动检查、数据库迁移和运维告警,不能只依赖重试。
multiGet 和 Pipeline;两个消费组直接消费原始 Topic,意味着标准化逻辑执行两次。这是用计算量换取故障隔离和低延迟的明确取舍。
降低额外成本的方法:
当流量规模继续扩大,且双标准化成本超过可接受范围时,可以重新评估标准化 Topic,但必须保证标准化生产层本身可水平扩展、有独立 Lag 监控,并通过压测证明不会重新成为共同瓶颈。
以下配置仅展示结构,Topic、消费组和地址必须根据环境管理,认证信息不得写入文档或代码仓库。
spring:
kafka:
consumer:
# 持久化链路复用原消费组位点。
group-id: ${PERSISTENCE_RAW_GROUP:<persistence-raw-group>}
max-poll-records: ${KAFKA_MAX_POLL_RECORDS:200}
properties:
allow.auto.create.topics: false
platform:
iot:
telemetry:
enabled: true
consumer-enabled: true
raw-topic: ${RAW_TELEMETRY_TOPIC:<raw-topic>}
# split:实时与持久化使用独立消费组直接消费原始Topic。
processing-mode: split
realtime-consumer-group: ${REALTIME_GROUP:<realtime-group>}
realtime-concurrency: 3
point-cache:
local-enabled: true
local-maximum-size: 50000
local-model-maximum-size: 10000
local-expire-after-access: 30m
redis-ttl: 2h
redis-ttl-jitter: 10m
redis-failure-cooldown: 30s
invalidation-channel: <cache-invalidation-channel>
sampling:
enabled: true
min-interval: 1s
max-interval: 1d
watermark-ttl: 7d
persistence:
max-points-per-write: 1000
replay-enabled: false
实际项目可能使用不同配置前缀,迁移到其他系统时应保持配置语义,不应照抄示例命名。
InfluxDB URL、组织、Bucket、Token,Kafka SASL/SSL 凭据,Redis 密码和配置中心凭据必须通过受控环境变量或密钥管理系统注入。
服务启动时建议输出一条安全摘要:
mode=split
rawTopic=<masked>
realtimeRawGroup=<masked>
persistenceRawGroup=<masked>
samplingEnabled=true
localCacheEnabled=true
不要输出 Token、密码、完整连接串或原始遥测载荷。
健康检查至少验证:
脱敏示例:
Realtime batch processed records=120 receivedPoints=840
coalescedPoints=390 published=390 costMs=18
字段含义:
receivedPoints:标准化后的合法点位数;coalescedPoints:批次内按设备和点位合并后的数量;published:实际发布到实时投影的数量;costMs:批次总处理时间。脱敏示例:
Persistence batch processed records=120 candidates=840
selected=140 written=140 skipped=698 outOfOrder=2
samplingCostMs=3 influxWriteCostMs=14 watermarkCommitCostMs=2 costMs=23
重点观察:
selected 应与当前点位间隔配置相符;written 应等于成功写入的选中点数;skipped 增加说明入库间隔正在生效,不代表数据异常;outOfOrder 持续增加需要检查设备时钟或历史补传;influxWriteCostMs 高于总处理预算时,应优先检查 InfluxDB 和网络;watermarkCommitCostMs 异常时检查 Redis,但不能让其掩盖 InfluxDB 结果。指标标签不得直接使用设备ID或点位编码,避免监控系统出现高基数。设备级问题通过日志、Trace 或受控诊断接口定位。
测试条件:
预期:
skipped 明显大于0。测试条件:
预期:
测试条件:
预期:
测试条件:
预期:
split。samplingIntervalSec = null 的点位每条合法消息均可入库。本次实践最重要的结论不是简单地“增加一个采样判断”,而是先把实时性和历史持久化的可靠性边界拆开:
对于工业物联网系统,实时和历史不是同一个消费语义。只有把二者拆成独立可观测、可失败、可恢复的链路,才能同时获得及时的页面体验和可控的时序存储成本。