
开头这里用一段引言不设标题这两年欧洲很多水务公司都在推进“智慧水务”项目但真正敢把实时调度和水质高并发分析放在同一个平台里做的并不多。我在奥斯陆参与的这套智能水利平台算是一个比较典型的工程例子既要根据实时水压、流量、蓄水池液位和天气预报去做水资源调度决策又要同时处理分布在全城各处水质传感器传来的高并发数据对氨氮、浊度、pH、余氯等指标做分钟级甚至秒级的分析。项目上线后这套平台在多次强降雨和用水高峰中扛住了压力也踩了不少坑。如果你正在做类似方向的系统——无论是水务、能源还是其他需要“实时采集高并发计算快速调度”的物联网平台这篇分享会把我们的架构设计思路、关键选型理由、以及那些只有真跑过才会知道的工程细节一次性讲透。文章里所有经验和数据都来自我个人的实际参与不是泛泛而谈的方案介绍。1. 项目背景为什么奥斯陆需要一套实时水利调度平台1.1 乍看是“监控大屏”实则是实时决策系统当初项目立项时甲方的表述很朴素我们要做一个能实时显示水质数据的大屏最好还能辅助调度员开关阀门、调节泵站。很多团队一听“大屏”就想着用WebSocket图表库把数据刷起来但我们深入调研后发现核心难点远不止“显示”。奥斯陆的供水网络有超过一千公里的管道水源来自多个湖泊和地下水井部分管网已经有几十年历史拓扑结构复杂。调度员需要根据水厂出水量、区域用水量、泵站压力、高峰期蓄水池余量来判断未来几个小时内是否需要调节阀门和泵机。这里面有明确的“调度”动作所以系统首先是一个实时决策系统其次才是监控系统。大屏只是实时特征的可视化窗口平台真正的价值在于毫秒级的数据到达、秒级的水质计算、分钟级的趋势预判以及发生异常时能不能自动给出调度建议。1.2 硬性指标实时不是口号而是可量化的协议项目在设计之初就定了几个硬性指标这里我直接列出来后面很多技术选型都跟这几个数字有关水质传感器数据从采集到进入内存计算引擎的端到端延迟不超过3秒。核心水质指标pH、浊度、余氯、氨氮的计算结果达到准实时刷新率默认5秒更新一次。平台需支持至少2万个传感器并发连接并且在峰值时每秒处理5万条以上的时序数据点。调度指令从确认下发到泵站/阀门的响应时间不超过2秒。历史数据需要支持无缝回放用于事后分析和模型训练。当时看了这几个指标团队里一个老工程师说了一句很到位的话这系统本质上是一个行业化的“高并发IM服务”只不过消息体是水质数据客户端是传感器而“聊天”是传感器和调度引擎之间的对话。后来我们还真借鉴了一些IM系统的设计思路比如长连接管理、消息队列削峰、Redis缓存做会话状态等。1.3 约束条件老旧管网、异构设备、环境敏感除了指标还有一堆现实约束。传感器来自四五个不同厂家协议不统一有的走Modbus RTU有的走LoRa有的走NB-IoT有的直接通过4G网关把JSON Post过来。数据质量也是大问题水下设备经常漂移偶尔还会上报负值的pH或者凭空消失的数据包。另外奥斯陆地处北欧冬季寒冷部分户外站点供电不稳定网络也有间歇性中断。这些约束意味着我们不能把系统设计成一个“完美网络环境下的玩具”——必须有完整的断线续传、数据缓存、质量清洗和自动补报机制。这也是为什么项目里我们把相当一部分精力放在接入层和数据处理链路上而不是只做漂亮的前端图表。2. 整体架构演进从“集中式入库”到“分层实时链路”2.1 第一版架构的教训一开始我们采用的是比较传统的方案传感器数据通过网关直接写入KafkaKafka下游接一个流处理程序写时序数据库然后应用层查询数据库刷新大屏。这套方案跑 demo 没问题但等到接入超过5000个传感器时问题全冒出来了。首先是Kafka - 流处理 - 时序数据库这条链路里流处理程序承担了太多任务协议解析、点位映射、异常检测、阈值判断、业务规则全揉在一起一旦某个传感器数据格式异常整个拓扑的重启会把延迟拖到十几秒。其次调度端的很多需求是“当前时刻的综合指标”例如某个水厂的整体出水质量得分、某个分区的余氯趋势斜率这些计算如果每次实时查询都去数据库聚合数据库压力非常大。当时我们用了开源的时序数据库在每秒2万点的写入下还能撑住但同时来十几个聚合查询时CPU和IO就开始打架查询延迟变得不可接受。最后我们团队内部达成了一个共识实时系统不能把所有逻辑都塞在“入库后查询”这一层必须把计算拆出去让数据在流动的过程中就完成大部分加工存储只是其中一个环节。2.2 正式方案三层解耦实时特征服务经过两轮重构我们最终采用了这样的架构接入与传输层负责设备接入、协议解析、断线缓存、数据清洗统一输出为标准化内部消息体。实时计算层由流处理引擎和内存计算引擎组成负责窗口统计、水质指标计算、异常检测、实时特征生成。服务与存储层对外提供查询API和WebSocket推送同时负责历史数据归档、冷热分离和与业务系统的数据同步。这里有一个关键设计——实时特征服务。我们把所有需要“低延迟读取”的派生指标单独维护在Redis和内存Grid中例如某个分区的平均压力、某水厂的实时出水流量、最近5分钟余氯变化斜率。计算层每隔几秒更新一次这些特征值而后端的调度决策服务、大屏可视化、告警服务都通过统一API读取特征值而不是各自去查数据库或重新计算。这个思路在后来支撑高并发查询时起了很大作用。大屏刷新、手机端App、调度台、外部数据共享所有这些访问都直接走特征服务数据源只有一份且始终新鲜。即使后端的时序数据库偶尔抖动大屏上的核心指标依然不受影响。2.3 选型理由与替代方案对比在流处理引擎选型上我们对比过Flink、Kafka Streams和Spark Streaming最终选了Flink。原因很简单我们需要毫秒级的事件时间窗口、精确一次语义以及能方便地与外部数据库做维表关联。Flink的Checkpoint机制在断网重连和故障恢复上非常可靠正好匹配我们传感器网络不稳定的特点。Kafka Streams虽然轻量但窗口和状态管理能力弱一些Spark Streaming的延迟是秒级甚至分钟级对于5秒更新一次的业务场景显得不够从容。在实时数仓同步工具上我们评估过几个开源增量同步方案。因为上游核心业务数据存在MySQL和PostgreSQL里我们希望把业务系统的订单位、设备台账等变更实时同步到分析型存储中最终选了基于Canal的增量同步链路再配合一款轻量的同步工具做结构映射。如果你不是特别依赖特定生态Canal DataX或者Debezium Kafka都是很成熟的组合。关键是不要自己写轮子去拉全量数据一定要基于binlog或WAL做增量同步否则数据延迟和压力都不可控。在数据库选型上我们用VictoriaMetrics作为时序存储而不是InfluxDB或TimescaleDB。原因之一是VictoriaMetrics的存储压缩率高二是在高并发写入下内存占用更稳定三是支持Prometheus协议方便和告警监控体系直接对接。如果你已经有Pushgateway或Prometheus技术栈选VictoriaMetrics会非常顺。如果团队对SQL特别依赖TimescaleDB也可以考虑但要注意它的写入放大问题。3. 实时数据采集与接入层高并发与设备协议并存3.1 设备接入协议栈先统一再分发接入层是整个平台最杂的部分但也是最不能乱的部分。我们做了一个“协议转换网关”的抽象层把不同厂家的协议全部转换成统一的内部数据模型。内部模型用Protobuf定义消息体里包含设备ID、点位编码、采样时间、接收时间、原值、量纲、质量标记、网关ID。这个设计让后续所有模块都不用关心数据来自什么设备只需要处理统一结构。协议接入主要分三类TCP长连接网关针对支持Modbus TCP或私有TCP协议的在线监测仪网关维护长连接主动拉取和被动接收都支持。MQTT接入服务针对采用MQTT协议的NB-IoT传感器网关作为MQTT Broker的客户端订阅主题再解析payload。HTTP回调接收服务针对一些4G/5G网关传感器按固定频率POST JSON数据。为了支持高并发连接我们用Netty实现TCP网关线程模型采用主从Reactor单个网关节点可以稳定维持1万个TCP长连接。MQTT接入用了EMQX集群保留会话和离线消息传感器断线重连后可以继续接收指令。3.2 数据预处理与质量标记原始数据必须做预处理否则脏数据会一路污染到调度决策。我们每个协议解析器后面都挂了一条小的清洗链格式校验字段是否齐全、类型是否正确、时间戳是否合理。范围检查根据点位定义的范围过滤物理不可能的值如负的浊度、pH14。跳变检查如果当前值与上一有效值变化超过物理上限先标记为“可疑”而不是直接丢弃。设备状态校验传感器心跳是否正常、网关是否处于校时状态如果是对应数据标记为“降级”。每个数据点会带上quality_flag字段取值包括GOOD、SUSPECT、BAD。后续计算引擎会自行决定是否将SUSPECT数据纳入统计。如果全部为BAD系统会生成一条设备离线告警而不是显示一个荒谬的值。3.3 高并发写入的缓冲设计接入层的计算压力主要在序列化和缓冲上。为了避免传感器上报的随机波动导致后端存储过载我们在接入层和计算层之间用了Kafka作为削峰缓冲。Kafka的分区键按设备ID哈希确保同一设备的数据进入同一分区从而保证单设备内的时间顺序。这里有个很重要的细节Kafka的分区数不能拍脑袋定。我们按目标吞吐5万点/秒、单分区写入能力1万点/秒来算预留一定余量设置了20个分区。如果分区太少消费并发上不去分区太多又会在故障恢复时拖慢均衡速度。建议每个分区对应的消费线程处理的数据量不要超过这个线程的实时计算能力否则会堆积。与高并发IM系统类似消息体设计时也要控制体积。很多传感器原始JSON到了接入层会被转成紧凑的二进制格式减少网络和存储开销。这不是为了炫技是真的能让Kafka吞吐提升不少。4. 内存计算层水质分析引擎与实时特征服务4.1 水质指标计算逻辑水质分析不是简单地“看数值超不超过阈值”它需要结合多个指标做综合评估。我们的水质分析引擎采用规则引擎轻量模型的方式。以“余氯”为例单个采样点的余氯值会被流处理任务放到一个5分钟的滑动窗口里计算平均值、最大值、最小值、标准差。当标准差大于设定值时说明数据抖动过大即使平均值合格也会触发“稳定性异常”。浊度则更多看趋势连续3个计算周期上升且斜率超过阈值系统判定为“恶化趋势”提前预警。针对氨氮我们使用卡尔曼滤波平滑原始值消除一部分传感器随机噪声再把平滑后的值送往规则引擎。这个处理很有效但要注意卡尔曼滤波的Q和R参数需要根据传感器类型单独调不能一套参数走天下。当时我们为了减少工作量按传感器型号分组调参实测下来误报率下降了40%。4.2 实时特征服务如何支撑调度实时特征服务是整个平台的“中路枢纽”。它由Flink作业实时计算并维护在内存Grid中同时定期把最新特征快照写入Redis缓存。举几个例子水厂出水特征出水量、出厂压力、单位能耗、水质综合评分。片区用水特征当前用水总量、预测未来2小时水量、供需缺口。水质异常事件特征最近1小时内异常次数、当前是否处于应急模式。调度决策服务在做方案推演时会读取这些特征结合管网水力模型做“what-if”分析。比如预测到某区域2小时后会缺水系统会先在虚拟环境里模拟“调高上游泵站转速”的效果看压力和流量是否超限如果通过再生成调度建议。这套逻辑如果用传统数据库查历史曲线来算几秒钟都跑不完一次推演而我们通过实时特征服务单次推演耗时稳定在300毫秒以内。4.3 异常检测规则与模型联动除了阈值规则我们还在流处理作业里部署了一个轻量的孤立森林模型用于检测多指标联合异常。例如在正常情况下pH和溶解氧之间存在一定的相关性单一指标都正常但相关关系突然断裂时往往意味着传感器故障或者水体受到污染。孤立森林可以捕捉这种非典型状态推理延迟在毫秒级。不过模型不能直接出告警否则会把运维人员淹没。我们设计了一个“可信度分级”机制模型输出异常得分如果得分不高仅写入异常事件表供追溯如果得分高且同时触发两条以上规则才会生成告警。这个机制上线后真正的人工介入次数降低了很多。5. 数据分发与存储Redis缓存设计、增量同步与历史归档5.1 Redis缓存让查询高并发不再难很多做高并发的人对Redis都不陌生但水务场景里Redis怎么用好跟纯互联网场景还是有些区别。我们主要用Redis做三件事第一是实时特征缓存。Flink每5秒计算一次特征后把结果写到Redis的Hash结构里key为特征类型加维度IDfield为具体指标。读侧服务统一走Redis命中率高而且天然支持集群分片。第二是设备状态缓存。为了支持设备列表的高频刷新我们维护了所有网关和传感器的最新状态在线、离线、数据延迟、电量等。这部分数据读多写少用Redis的String过期时间简单高效。第三是调度指令状态缓存。调度指令下发是一个异步过程指令发往设备后设备可能延迟确认。我们用Redis存储指令的当前状态并通过Pub/Sub通知相关模块状态变化。在缓存一致性的处理上我们的经验是“先更新存储再更新缓存最终通过TTL兜底”。直接删缓存而不是更新因为特征数据是整块重算的优势明显。TTL一般设置成计算周期的2倍避免缓存和实时计算不同步。5.2 增量同步机制从业务库到数仓平台不仅处理传感器实时数据还要关联业务系统里的设备台账、订单位、维护工单等数据。这些数据存在MySQL里且需要以准实时的方式同步到分析型存储和缓存中。选型时我们调研过几款MySQL增量同步工具。那段时间正好社区里有人推荐基于Canal的同步方式我们也研究了如Debezium之类的方案。总结下来Debezium的生态好、可扩展性强适合Kafka系而如果团队以Java为主Canal更轻、更容易定制解析逻辑。我们最终用了Canal 自研的Sink Connector把binlog解析出的变更事件写入Kafka再由Flink作业做维表关联和缓存刷新。这个链路有个很重要的点增量同步必须保证幂等。因为同步过程中可能重复消费binlog事件如果Sink不是幂等的对账就会出现问题。我们在写入分析库时采用“主键更新时间”的upsert方式保证同一行数据重复写入时结果一致。5.3 冷热分离与分区策略时序数据会无限增长不可能把所有数据都放在热存储里。我们的策略是热数据保留在VictoriaMetrics中最近30天按时间分片超过30天的数据转存到对象存储Parquet文件用于历史回放和模型训练。分区策略上我们按天创建数据分区同时按标签做二级分片。写入时通过设备所属区域路由例如“Oslo-Center”和“Oslo-West”的数据落到不同的底层分片这样区域维度的聚合查询不会跨片扫描全部数据。查询优化时我们尽量利用时间范围过滤因为大多数实时分析都聚焦最近5分钟或最近1小时很少需要全量扫描。对象存储归档我们布置了一个定时任务每天凌晨把前一天的时序数据导出为Parquet并自动清理热库中的过期分片。导出时按设备ID和点位编码聚合成大文件避免小文件过多导致查询性能下降。6. 高并发场景下的性能压测与排查实录6.1 我们是怎么做压测的项目上线前我们搭了一套与生产环境配置相当的压测环境用JMeter和InfluxDB数据生成器模拟2万个传感器上报。压测计划分三档正常峰值、1.5倍过载、2倍极限。具体做法是先跑一段基线流量观察各环节CPU、内存、GC、网络带宽。逐步提高并发连接数和每秒消息数每档持续30分钟记录延迟分位数。在2倍极限档位下人为杀掉一个Flink TaskManager节点观察恢复时间。这里有个经验压测不能只压接入层一定要全链路压测包括Kafka、Flink、Redis、数据库、WebSocket推送服务否则很容易漏掉瓶颈。我们的第一轮压测就漏了WebSocket推送服务结果接入层和数据计算都扛住了大屏前端因为推送量过大出现卡顿。6.2 三个典型性能问题与根因问题一Redis连接池被打满现象是高峰期部分查询请求超时Redis监控显示连接数接近最大限制。排查发现是我们的查询服务用了同步客户端每个线程在读到结果前会占用一个连接而服务线程池开得很大直接把连接池打爆。解决办法一是改用Lettuce异步客户端减少线程阻塞二是调整连接池最大连接数和最大等待时间三是把实时特征查询进一步收敛到几个聚合接口避免前端低频发请求。问题二Flink作业反压导致端到端延迟升高现象是所有数据都能收到但水质计算更新频率从5秒变成了15秒。排查Flink UI看到某个算子反压百分百。最终定位到问题出在“维度关联”算子它需要访问MySQL里的设备台账但为了省事我们直接用一个外部查询Client每次去查库MySQL扛不住。解决办法把设备台账加载到Flink的广播状态中定时更新比如每分钟刷新一次让维度关联走内存彻底消除外部依赖。问题三GC导致查询长尾现象是大部分查询在30毫秒内返回但时不时出现500毫秒以上的长尾。排查JVM日志发现堆内存在大批短生命周期对象。最终定位是创建了大量冗余的日志对象和序列化临时对象。解决办法调整日志级别去掉debug日志对热点路径复用对象池把部分对象改为栈上分配。优化后P99从800毫秒降到了60毫秒。6.3 优化后能抗住的流量经过优化最终生产环境在2万个传感器并发连接下每秒写入约5.5万时序数据点端到端数据到达延迟小于1.5秒目标3秒水质特征更新周期稳定在5秒调度指令下发响应时间在1秒左右。大屏WebSocket推送支持20万在线客户端时依然流畅。这些数字印证了一件事只要分层合理高并发实时系统不是靠运气而是靠每个环节的可控设计。7. 工程化落地中的避坑指南与实际心得7.1 最容易忽略的时间同步问题在海量设备接入时设备时间不一致会引发大问题。比如某传感器上报的时间比服务器早了10分钟如果不加处理按事件时间计算的窗口统计会乱套实时特征可能出现“跳到未来”的假象。我们的方案是所有设备上报必须携带设备时间和网关接收时间进入平台后统一转换为UTC时间戳并以网关接收时间为准判断延迟。如果设备时间与网关时间偏差超过5分钟该设备的数据自动标记为“时间漂移”不参与实时计算直到校时完成。同时所有流处理的时间窗口都基于事件时间配合Flink的Watermark机制处理乱序数据。7.2 千万别把所有数据都当成“实时”实时系统里最容易犯的错就是把所有数据都追求秒级。其实有很多数据是“准实时”就够了的。比如设备台账更新、维护工单状态、用户账单信息这些数据延迟几分钟根本不影响调度。一开始我们想把所有同步都做实时结果Canal同步链路和接口调用异常复杂还经常因为非核心数据问题影响核心链路。后来我们在架构上做了区分核心水质和压力数据走高保真实时链路设备管理、工单等业务数据走“准实时定时刷新”链路。这样既保证了关键路径的稳定又降低了系统复杂度。工程领域有一句话最快不是最优可控才是王道。特别是实时调度平台稳定性和可预期性远比绝对延迟重要。7.3 运维监控与告警设计实时平台自身的可观测性非常关键。我们为所有关键组件都接入了指标采集包括Kafka消费延迟、Flink反压指标、Redis命中率、接入层连接数、消息积压数等。告警规则分成三级提示级不打扰人、警告级发送到值班群、严重级电话通知自动降级策略。特别是消息积压告警。我们在Kafka上给每个主题都设置了消费延迟阈值一旦超过阈值就自动扩容Flink作业并行度。这里要强调自动扩容的逻辑必须经过充分测试否则扩容瞬间的状态恢复反而会造成二次冲击。我们只对无状态作业做自动扩容有状态作业还是人工干预更安全。监控大屏自身也有个容易忽略的点要盯着“数据新鲜度”而不是只盯着页面有没有报错。如果数据流一切正常但某个图表已经5分钟没更新了这往往说明链路中有静默故障。我们增加了一个“新鲜度”看板专门展示每条链路从采集端到展示端的时间差一旦时间差超过设定值就高亮告警。写在最后的一个小经验如果你也在设计类似的实时平台我最想强调的一点是实时系统的难点不在某个单项技术而在于把延迟预算分摊到每个环节。要明确哪些路径必须快、哪些路径可以慢哪些数据必须准、哪些可以容忍误差。别把所有问题都甩给“Flink”或者“Redis”解决很多瓶颈是架构层面提前埋下的。另外团队里最好有一个人始终盯住“端到端体验”——他不是某个组件的负责人而是所有数据路径的owner。我们项目里这个角色就是我自己这让我能及时发现“虽然每个模块都正常工作但整个链路已经达不到验收标准”的尴尬局面。最后再分享一个实用小技巧图像或地图类可视化页面的实时刷新频率一定要单独做动态调节。当地图上有大量点位密集出现时5秒刷新会卡成PPT这时候如果能把点位聚合后推送体验会好很多。这可以说是所有实时可视化平台都值得留意的细节。真希望当年第一版上线前就明白这些能少走不少弯路。