我之前做过不少水利相关的数据项目,但像这次这样把高并发实时调度和水质分析揉在同一个平台里的,还真是头一回。整个项目从需求梳理到落地,踩了不少坑,也沉淀了些可复用的思路。今天就把奥斯陆这个智能水利场景下的工程设计实践,从头到尾拆开聊聊。
1. 项目背景与整体设计思路
1.1 业务痛点与核心需求拆解
奥斯陆的城市供水系统覆盖了水源地、多个水处理厂、加压泵站和上千个管网监测点。水务公司原有的架构是典型的“采集-入库-展示”模式:传感器数据定时上报,存进关系型数据库,再由后端服务按需查询。这套模式在数据量小、实时性要求低的场景下能用,但一旦设备数量上去、上报频率加密,问题就全暴露出来了。
项目启动时梳理出的核心痛点很清晰:
- 调度决策依赖的用水量预测数据是“偏昨天的”,无法支撑当天的高频调度。
- 水质监测还是小时级抽检模式,突发污染根本来不及预警。
- 泵站和阀门控制依赖人工电话沟通,一条调度指令从下达到执行往往要几十分钟。
- 上千个传感器高频上报后,原来的单体数据库频繁出现连接数打满、查询超时。
所以这个平台的建设目标,我拆成了四个层次:
- 第一层是数据接入层,要能扛住高并发写入,不能丢数据。
- 第二层是计算层,实时做用水预测、水质异常识别、调度策略计算。
- 第三层是控制层,把计算结果转化为泵站、阀门的调节指令,并保障可靠下发。
- 第四层是业务层,给调度员看板和大屏,让人能看懂、能干预、能追溯。
一句话总结:这套系统要替代原来的人工经验和定时报表,变成24小时不停歇的“调度大脑”。
1.2 总体架构选型:从定时任务到实时流式处理
架构选型时,我们反复对比了“加大单体服务器”和“拆成流式架构”两条路线。最简单的做法是把原来自研的轮询采集改成消息队列加批处理,但算下来高峰期每秒几千条写入,批处理窗口要是卡住,调度延迟直线上升。我们当时拍板的原则是:核心链路必须端到端秒级延迟,非核心环节允许分钟级。
最终确定的总体架构按数据流向分成了五层:
- 感知层:水质传感器、流量计、压力计、水位计、气象站,统一通过边缘网关汇聚。
- 接入层:网关上报的数据统一进入Kafka消息集群,这里作为全链路缓冲和削峰填谷的“蓄水池”。
- 计算层:Flink负责实时流处理,包括数据清洗、指标计算、用水预测、异常识别;Spark批处理负责离线模型训练和历史数据重算。
- 存储层:时序数据库存原始监测值和实时指标,关系型数据库存设备档案与调度指令,对象存储存原始日志和模型文件。
- 应用层:调度控制台、水质预警、可视化大屏、报表系统。
选型时我们没有引入太多新玩意,尽量选团队熟悉、社区活跃的组件。Kafka扛高并发写入很成熟,Flink在流计算领域基本是事实标准,时序库用的是开源方案部署在私有云上。整体思路是“用最稳妥的方案拼出最需要的效果”。
这套架构最舒服的一点是:每一层都可以独立扩容。接入层压力大了加Kafka分区,计算层吞吐不够加Flink并行度,存储层容量紧张了加节点。不像原来那个单体库,压力一大就只能整体迁移。
2. 高并发数据接入与消息链路设计
2.1 边缘网关接入方案的选型
整个平台最基础也最容易翻车的地方,就是传感器数据怎么稳定、实时地到达平台。奥斯陆这个场景里,传感器分布在几十公里长的管网上,有的在偏僻的郊外泵站,有的在市区地下管廊,网络环境差别很大。
我们的方案分两个层面:
硬件侧,统一部署边缘网关,网关接入传感器后做三件事:一是数据缓存,断网时本地暂存,恢复后补传;二是协议转换,不同厂商的传感器协议统一转成MQTT协议;三是初步过滤,明显超出量程的死值直接丢弃,减少无效上报。选网关时核心指标是宽温范围和断网续传能力,挪威这边冬季气温低,设备稳定性必须过硬。
软件侧,采集服务部署成集群,通过域名接入网关上报。采集服务只做一件事:把收到的数据包装成统一格式转发到Kafka,不落库、不重计算。这个服务把状态全部外置到Redis里,自身做成无状态,所以可以随意水平扩容。
这里有个踩过的坑值得说:一开始我们把协议解析放在采集服务里,结果每个厂商的传感器都要改代码,上线后维护成本很高。后来我们把协议解析下沉到边缘网关,平台上只接收标准化的JSON数据,新增传感器型号只需要在网关上做配置,平台层完全不用动。
2.2 用Kafka做流量缓冲,解决峰值写入压力
传感器数据的写入特点是平时平稳、早晚用水高峰时流量猛增,同时突发情况下(比如爆管、水质报警)会有大量事件同时上报。如果让这份压力直接打到数据库,再好的数据库也会被打满连接数。
所以Kafka在这个架构里的角色是标准的“削峰填谷”。我们在Kafka里按业务类型建了三个核心主题,数据从网关到Kafka的主题,就是最简单的生产消费模型,生产端是采集服务,消费端是Flink计算任务。
分区数的设计是根据消费端并行度和吞吐要求反推的。我们按传感器类型做分区路由,水质数据按设备ID哈希分到不同分区,调度事件按时间顺序保证同一设备的事件有序。这个有序性当时花了不少功夫测试,因为后续计算需要保证同一监测点的数据严格按照时间顺序被处理。
Kafka集群的配置上,我们开了三副本保证高可用。节点挂了自动切换leader,生产端配置重试机制,消费端记录offset。这套组合拳打下来,线上跑了半年多,没有丢过一条真实数据。
另外给一个实操建议:网关上报频率不要拍脑袋定。我们验证过,水质常规指标15分钟上报一次足够,但管网压力在调度期间要1分钟一次。频率定高了,Kafka和存储的成本直线上升;频率定低了,调度反应不过来。这个频次要和业务方逐类确认,用最小成本满足最短业务闭环时长。
2.3 数据语义校验与异常过滤
数据进入平台不等于数据能用。高并发场景下,传感器故障、网络抖动、电磁干扰导致的异常数据掺杂在正常数据流里,如果直接把这些数据送进计算引擎,计算出的调度策略全会偏掉。
我们在Flink计算链路里,实现了一个独立的数据质量算子,专门做三件事:
- 完整性校验:上报数据缺少必要字段(比如没有设备编号、没有时间戳)直接丢弃并记录日志。
- 范围校验:超出传感器量程的数据直接丢弃。比如pH值正常范围是0到14,上报个20就明显不可能。
- 突变校验:对比该监测点最近N条数据的滑动平均值,超过阈值3倍以上的波动,标记为疑似异常,进入人工确认队列,不参与实时调度计算。
这个突变校验逻辑帮我们抓住过不少问题。有一次某泵站的压力传感器出现间歇性故障,每隔一段时间上报一个异常高值,原来的规则引擎会直接触发报警,而我们的算法把它识别成了疑似异常,调度员在确认界面上看到并核验后,避免了误触发降压操作。这个小功能对业务来说很不起眼,但实际价值非常大。
3. 实时水资源调度模型的工程化实现
3.1 短期用水预测与调度决策
水资源调度的核心是“提前判断用水趋势,预先安排泵站和水库的出力”。如果等用水量已经上来了再反应,管网压力已经出现波动了。
我们的调度模型分两层:日级规划和实时调整。
日级规划在每天早上跑一次,基于历史用水数据、天气预报(温度、降水)、是否为工作日/节假日等特征,用时序模型预测当天24小时的用水曲线。预测结果作为当天调度的基准线,输出到调度平台展示。
实时调整是每5分钟执行一次:根据最近1小时的实测用水数据,用轻量级回归模型修正“日级预测偏差”,动态调整未来15分钟的泵站目标压力。之所以用轻量级模型,是因为要在Flink里跑,计算延迟必须控制在秒级。
模型训练用历史数据在Spark里批量完成,训练好的模型导出为文件,定期滚动上线。实时计算时Flink加载模型文件,把特征向量推进去得出预测结果。训练和推理分离,线上计算非常轻,这是工程上最常见也最稳妥的实践方式。
预测效果上,整体平均绝对百分比误差控制在10%以内。但必须承认,节假日和极端天气下的预测偏差仍然比较大,这块不是纯靠模型能解决的,我们结合了人工预案:遇到特殊天气预警,调度员可以一键切换为“特殊天气调度模式”,调用另一套针对性的控制参数。
3.2 泵站与阀门控制指令的下发链路
调度模型算出的结果要执行,必须通过泵站和阀门的远程控制。这个环节最容易出问题,因为控制指令是反向链路,一旦出错,影响是实打实的物理世界。
我把控制指令的下发链路设计成三步:
第一步,调度计算模块产出指令后,先写入指令表中,状态为“待确认”。这里不直接下发,是为了留一层人工审核的余地(可以配置为自动下发模式)。
第二步,指令下发服务从指令表里轮询待确认指令,通过加密通道下发到对应泵站或阀门的边缘控制器。边缘控制器执行后将结果返回。
第三步,指令执行结果回流到平台,更新指令状态,同时调度大屏上实时展示“指令下发成功、等待反馈、执行完成”等状态。
这个链路看上去简单,但细节问题很多。比如指令下发需要幂等控制,同一指令不能因为网络重试而被执行两次,否则泵站会来回启停。我们在指令里加了全局唯一的指令ID,控制器侧按指令ID去重,重复的指令直接忽略并返回已处理状态。这就是典型的工程细节救命的场景。
另外,控制链路的监控比数据链路更重要,我们对每条指令的每个状态流转都做了日志记录和耗时统计,超过10秒没有反馈的自动转人工处理。因为泵站控制一旦出现“指令已下发但无反馈”的模糊状态,调度员必须最快速度介入,否则可能影响供水压力。
3.3 调度效果评估与闭环反馈
调度模型不能只算不管效果,否则模型漂移了没人知道。我们在平台里建了调度效果评估模块,每次实时调整执行后,持续跟踪后续15分钟内管网压力实际变化,和预期效果做对比,计算偏差指标。
这个指标会每天汇总一次,按泵站、按时段分析调度策略的准确性,形成一份自动化日报。月度会议上,水务调度团队会拿这份报表复盘,找出一段时间内频繁出现偏差的区域,针对性调参或修改模型逻辑。
这实际上形成了一个完整的闭环:预测-决策-执行-评估-再优化。当初定这个机制时,模型团队是有顾虑的,怕效果评估结果不好看。但实际上线后,评估机制逼着模型持续迭代,半年时间调度的精准度提升非常明显,水务公司也认可了这个机制的价值。
我个人认为,任何一个“带决策性质”的平台,如果没有效果评估闭环,都只能算是一个“高级看板”,谈不上真正的智能调度。
4. 水质数据分析平台:从清洗到可视化
4.1 水质指标存储与数据湖分层
水质数据的分析相比调度数据,最大的难点是数据维度多、关联关系复杂。一个监测点的水质数据包括余氯、pH值、浊度、溶解氧、电导率、温度等多达十几个指标,每个指标的数据特性都不一样,对存储和计算的要求也不同。
我们的存储方案是“时序库+对象存储”双轨制:
- 时序数据库存储近3个月的原始监测数据和实时计算指标,支撑快速查询和实时预警。
- 对象存储(兼容S3协议)存储全量原始数据,按日期和设备ID做分区。三个月前的数据自动从时序库转存到对象存储,报表查询走离线Spark任务扫描对象存储。
这套双轨制解决的核心问题是成本。时序库的资源开销和存储量成正比,如果无限期保留全量数据,成本压力会很大。对象存储便宜很多,虽然查询慢,但历史趋势分析并不需要毫秒级响应。
数据分层上,我们把数据湖划分为多级:
- 落地层:原始数据的明文备份,未经过任何加工。
- 清洗层:去除了异常值和空值,补全设备维度的元数据。
- 业务层:按业务口径预聚合后生成的主题数据,比如“某水厂进水口每小时平均浊度”。
- 应用层:直接服务应用和报表的高性能数据宽表。
每一层之间是单向依赖关系,不允许跨层调用。这条规则极大简化了数据血缘管理,出问题时能快速定位是哪个环节出错。
4.2 高吞吐量实时分析与预警服务
水质预警是平台里对延迟要求最高的模块。从传感器采集到异常事件推送,我们当时定的目标是10秒内完成,因为余氯异常这类问题越早发现损失越小。
实时分析链路的做法是:Flink消费Kafka水质主题数据,做窗口聚合和规则判断。规则引擎支持两类规则:
- 阈值规则:某个指标连续N次超过设定阈值则触发预警。
- 趋势规则:某个指标在时间窗口内持续上升且上升斜率超过设定值,即使尚未超阈值,也触发预警。
这个趋势规则很实用。比如浊度缓慢上升但还没到超限值,阈值规则不触发,等超限了可能已经晚了。趋势规则在浊度持续上升时就提前预警,调度人员来得及提前采取措施。
预警产生后,通过内部消息推送服务分发到短信、App推送和调度大屏三个渠道。消息推送服务用了独立的消息队列,避免预警高峰时阻塞核心调度链路。
另外我们给预警配置了自动降噪机制:同一监测点同一指标在10分钟内重复触发同一规则,会自动合并为一条预警,避免短信轰炸。否则管网压力波动时,一个点可能会一分钟发十条预警,调度员直接崩溃。
4.3 可视化大屏与报表体系
可视化部分的设计目标很明确:让调度员一眼看清水、压、质、量四个维度的实时状态。大屏上是七个核心模块:管网总览、区域压力分布、泵站运行状态、水质实时指标、预警滚动列表、调度指令执行追踪,以及今日用水曲线对照预测值。
大屏的技术栈上,我们直接用了前端团队熟悉的开源图表库加GIS引擎,底层接时序库提供的实时数据接口。数据刷新策略是WebSocket推送,服务端有数据变更才推,不走定时轮询,避免不必要的资源开销。
水务公司的管理层还需要日报、周报和月报,这类固定的报表我们直接用定时任务生成PDF,推送到邮箱和内部办公平台。对报表系统,我认为最重要的一点是口径要稳定。同一个指标(比如“日均供水量”),大屏上数字和日报上数字不一致,会被业务方反复质疑数据准确性,影响整个平台的信任度。所以指标口径在平台上统一管理,无论哪个入口展示,都走同一个查询服务。
5. 高并发治理与性能调优实录
5.1 突发流量下的限流与降级策略
这个项目虽然不像电商秒杀那样有超级流量高峰,但城市用水存在明显的“早晚高峰”,在这些时段里并发写入量和查询量都比较大,而且台风、暴雨等极端天气下,传感器上报频率和调度服务调用量也会猛增。一开始我们没做防护,结果就有一次因为一个外部接口调用链路过长,导致核心调度模块响应变慢,最后连锁影响到指令下发。
为了保障核心调度链路不受非核心服务影响,我们做了三层保护:
第一层是接口限流。接入层和业务接口统一做限流配置,针对不同接口设置不同的QPS上限,超过上限的请求直接返回“系统繁忙”的提示,而不是把压力传导到下游。
第二层是服务降级。非核心服务(比如报表服务、历史数据查询)在资源紧张时自动降级,优先保障实时调度链路。降级是预先配置好阈值的,比如当系统负载持续超过某个水平时,报表服务自动挂起,等负载降下来再恢复。
第三层是线程池隔离。核心服务用独立的线程池,非核心服务不能用完线程池里的线程。即使某个非核心服务被拖垮,核心服务的线程池也不受影响。
这套策略我当时是借鉴了业界做高并发微服务治理的常见思路。你不一定需要引入很重的框架,但“限流-降级-隔离”这三个手段必须有,至于用什么组件实现,反而不是最重要的事。
5.2 数据倾斜与热点问题的处理
高并发流式计算里,数据倾斜是个经典问题。我们这边表现最明显的场景是:某个人口密集区域的监测点数量特别多,按设备ID哈希分区后,这些数据都进了同一个Flink子任务,导致那个子任务的负载远高于其他子任务,出现“木桶效应”。
解决这个问题我们用了两步:
第一步是扩展分区维度。不再单纯按设备ID分区,而是按“区域+设备ID”的组合维度分区,这样数据分布更均匀。
第二步是加一层内部重分区。Flink任务内部分成两个阶段,第一阶段按原始维度处理单个设备的数据(比如指标清洗),第二阶段按更细的维度重新分区后做聚合计算(比如区域汇总),避免单点数据集中导致的热点问题。
另外还遇到过计算热点的场景:某个水厂进水口的流量数据量特别大,每次实时聚合计算都集中在这一个点上。我们的处理方式是对热数据做预聚合:在上游先把这一路数据聚合成更粗粒度的结果,下游计算再拿预聚合结果做二次汇总,计算压力直接降了一个量级。
5.3 资源成本优化与集群水位管理
流式平台稳定跑起来之后,资源成本优化就成了头等大事。Flink和Spark常年占着一批机器,尤其是Flink的常驻任务,任务多了之后资源浪费非常明显。
我们做了两项优化:
第一项是动态资源调整。根据业务高峰低谷的规律,给Flink任务配置动态扩缩容策略。用水早晚高峰时增加并行度,夜间低谷时缩减并行度,配合Kubernetes的Pod数量自动调整,实测成本下降约30%。
第二项是压缩存储成本。时序库的定期归档我们已经做了,后来又加了压缩策略,历史数据按更高压缩比存储,查询性能虽有小幅下降,但可以接受。
成本优化这个事,不能等老板提才做。技术上不难,关键是平时要统计好资源利用率数据,摸清每个任务到底消耗了多少资源,值不值得优化。这也是工程经验的一部分。
6. 常见问题与排查技巧实录
6.1 时间序列数据的乱序问题
实时场景下数据传输难免乱序:传感器A上报了10:00的数据,但因为网络原因,传感器B的9:58分数据才到。如果计算引擎按到达顺序处理,聚合窗口里就会出现“时间倒流”,结果混乱。
我们的解法是:Flink的窗口算子配置了事件时间语义和一定的乱序容忍度。事件时间以传感器采集时间为准,不是以平台接收时间为准。同时容忍度设置为30秒,超过30秒的迟到数据落到侧输出流,进入修正任务或丢弃该窗口的已计算结果,根据业务需求决定是否重算。
这个事件时间语义一开始我们没用好,直接用了处理时间,结果早晚高峰时数据乱序严重,用水预测频繁跳变。后来花了整整两天排查才定位到是时间语义的问题。换掉之后,结果平滑了很多。
6.2 传感器误报与数据质量追溯
水质传感器在运行一段时间后容易漂移,产生系统性偏差,这是硬件无法完全避免的。有一次一个监测点的余氯数据一直偏低,但没有触发阈值规则,因为偏差不大。直到我们的人工复核抽样才发现这个设备已经漂移了一段时间。
这个问题靠规则引擎很难识别,我们后来增加了在线校验机制:
平台会定期用同一区域内相邻监测点的数据做交叉比对,如果某个点的指标长期偏离区域均值达到一定幅度,自动打上“数据可疑”的标签,提示运维人员去现场校验传感器。这个交叉校验的逻辑,其实就是利用数据自身的冗余来发现单点异常,在传感器故障频发的区域,效果相当明显。
6.3 上下游链路压力不匹配
系统规模上来后,Flink处理能力扩容了,但下游存储集群和消息队列没跟上,结果Flink的吞吐量上去了,下游写入变慢,数据开始积压。
排查路径是先看Kafka各分区的消费延迟:如果延迟持续上涨,说明消费端(Flink)处理慢了;如果消费端很快而存储写入慢,说明存储成为瓶颈。我们用这套方法定位过一次时序库写入性能下降,原因是一个索引建得不合理,修复索引后写入速度立竿见影。
针对这类链路压力不匹配问题,我们建立了常态化的容量水位监控:每一条核心链路的Kafka消费延迟、Flink背压情况、存储写入延迟都做了监控大屏和阈值告警。哪一节链路水位异常,一眼就能看到。
6.4 实用的排查工具与调试技巧
最后分享几个工程排查的小技巧,虽然不起眼但很实用:
- Flink Web UI的Backpressure面板是定位任务瓶颈的第一入口,哪条链路背压高,问题基本就在哪。
- Kafka消费延迟建议用通配方式绘制趋势曲线,比只看当前值更能判断问题的趋势是持续恶化还是已经缓解。
- 排查乱序或丢数据问题时,抽样打印一条数据从网关到存储全链路的时间戳,会非常直观地暴露问题环节。
- 上线新规则或新模型前,先跑影子模式,同时运行新旧两套逻辑,只记录对比结果不下发真实指令,这是最稳妥的灰度方式。
我个人在实际操作中的体会是:在这样一套“监测+调度+分析”一体的平台里,最难的不是某个技术难点攻克不了,而是要对全链路的各个环节都有足够清晰的把握。数据怎么来、怎么算、怎么用、怎么反馈,哪一环出了问题,整个系统都会给出连锁反应。所以设计时一定要留出足够多的观测点,把监控和日志做好,比增加多少新功能都重要。
这套平台上线后,调度人员从原来的人工经验判断逐渐转为与系统协同工作,水质异常的响应时间从小时级缩短到分钟级。工程上走通之后,后续还有很多可以扩展的方向,比如接入更多降水与融雪数据来提升预测精度,或者把同类的调度能力复制到其他城市,这都取决于第一版的基础是否扎实。