Spring+大数据疫情监控系统架构与实现
·
1. 项目背景与核心价值
这个基于Spring+大数据的疫情监控系统是我去年指导的一个本科毕设项目,后来在实际应用中不断迭代完善。当时学生面临几个典型痛点:一是缺乏真实业务场景的毕设选题,二是对大数据处理流程不熟悉,三是远程协作开发经验不足。这个系统恰好能同时解决这三个问题——它既有公共卫生领域的现实意义,又融合了企业级开发中常见的Spring框架与大数据技术栈,还通过Git+Docker实现了完整的远程开发调试流程。
从技术角度看,该系统最核心的价值在于实现了三个"实时":
- 数据采集实时化:通过分布式爬虫集群抓取各地卫健委公开数据,延迟控制在3分钟以内
- 处理流程实时化:采用Lambda架构同时满足批处理和流式计算需求
- 可视化实时更新:基于WebSocket推送使前端大屏无需刷新即可同步最新数据
实际部署时发现,当数据源超过20个地区时,建议采用Kafka作为消息队列缓冲,否则直接写入HDFS会导致小文件问题
2. 技术架构解析
2.1 整体架构设计
系统采用经典的三层架构,但针对疫情数据特点做了特殊优化:
[数据源] -> [Flume+Kafka] -> [Spark Streaming]
-> [HBase] -> [Spring Boot] -> [ECharts]
↗ [Spark SQL] ↘
[MySQL] [Redis]
关键设计考量:
- 双存储引擎:HBase存储原始轨迹数据(满足CDC需求),MySQL存储聚合结果(支持复杂查询)
- 混合计算模式:Spark Streaming处理实时指标,夜间用Spark SQL跑批处理修正数据
- 缓存策略:Redis不仅缓存热点数据,还存储地理围栏计算中间结果
2.2 核心技术选型对比
| 技术点 | 候选方案 | 最终选择 | 选择理由 |
|---|---|---|---|
| 流处理引擎 | Flink vs Spark | Spark Streaming | 团队已有Spark经验,学习成本低 |
| 时序数据库 | InfluxDB vs HBase | HBase | 需要与历史数据联合查询 |
| 地理计算 | PostGIS vs JavaTop | JavaTop | 轻量级,满足省市级精度需求 |
| 前端可视化 | D3.js vs ECharts | ECharts | 疫情地图模板丰富,开发效率高 |
3. 关键实现细节
3.1 疫情数据采集模块
数据源处理采用分级策略:
- 省级数据:直接调用政务API(需处理HTTPS证书)
- 市级数据:Python爬虫集群(使用代理IP轮询)
- 民间数据:NLP清洗社交媒体文本(准确率约82%)
// 示例:多线程爬虫核心逻辑
@Async("crawlerThreadPool")
public void crawlData(DataSource source) {
while (true) {
try {
String html = Jsoup.connect(source.getUrl())
.timeout(10000)
.proxy(getRandomProxy())
.execute().body();
// 解析逻辑...
kafkaTemplate.send("raw-data", message);
Thread.sleep(source.getInterval() * 1000);
} catch (Exception e) {
log.error("爬取异常", e);
}
}
}
踩坑记录:某直辖市卫健委页面改版导致连续3天数据缺失,后来增加了XPath多套解析方案自动切换机制
3.2 实时计算模块
Spark Streaming处理流程优化点:
- 使用mapWithState替代updateStateByKey(内存消耗降低40%)
- 对确诊数等关键指标采用HyperLogLog去重计数
- 地理围栏计算采用R树空间索引
// 关键指标计算示例
kafkaStream.map(msg => {
val city = msg.geoHash.substring(0, 5) // 使用GeoHash分级
(city, msg.confirmed)
}).reduceByKeyAndWindow(
_ + _, _ - _,
Minutes(15), Seconds(30)
)
4. 远程开发实践
4.1 环境标准化方案
通过Docker Compose统一开发环境:
version: '3'
services:
hadoop:
image: apache/hadoop:3.3.6
ports:
- "9870:9870"
spark:
image: bitnami/spark:3.4
depends_on:
- hadoop
debug:
image: openjdk:11-jdk
volumes:
- .:/workspace
ports:
- "5005:5005" # 远程调试端口
4.2 典型问题排查表
| 现象 | 可能原因 | 解决方案 |
|---|---|---|
| Kafka消费延迟增大 | Spark批次间隔设置不合理 | 调整spark.streaming.blockInterval |
| 地图显示坐标漂移 | 坐标系转换未统一 | 所有数据强制转GCJ-02 |
| 定时任务重复执行 | 服务器时区不一致 | Docker统一设置Asia/Shanghai时区 |
| 内存泄漏 | HBase连接未关闭 | 使用try-with-resources重构代码 |
5. 项目扩展方向
在实际应用中我们还尝试了这些增强方案:
- 智能预警:通过LSTM预测未来7天传播趋势(需调整损失函数应对稀疏数据)
- 物资调度:结合高德API计算最优配送路径(注意避开封控区)
- 舆情分析:使用HanLP提取微博关键词(需处理大量网络用语)
部署建议:
- 生产环境建议至少3节点集群
- 每日凌晨执行HBase压缩(major_compact)
- 对公开API增加限流保护(Guava RateLimiter)
这个项目最让我意外的收获是:简单的技术组合也能产生实用价值。有家社区医院仅用2台旧服务器就部署了简化版,在2022年底的疫情高峰中发挥了重要作用。如果重新设计,我会在数据采集层加入更多验证机制,毕竟脏数据对实时系统的影响是指数级放大的。
更多推荐
所有评论(0)