结构健康监测仿真-主题044-结构健康监测中的云计算与大数据技术

目录

  1. 引言
  2. 云计算基础
  3. 大数据技术架构
  4. 数据存储与管理
  5. 数据处理与分析
  6. 机器学习与人工智能
  7. 可视化与交互
  8. 安全与隐私保护
  9. 工程应用案例
  10. Python仿真实现
  11. 总结与展望

在这里插入图片描述
在这里插入图片描述

引言

随着结构健康监测系统的规模不断扩大和数据量的爆炸式增长,传统的本地数据处理和存储方式已无法满足现代监测需求。云计算与大数据技术的引入,为结构健康监测提供了强大的数据存储、处理和分析能力,使得海量监测数据的高效管理和深度挖掘成为可能。

1.1 云计算与大数据的发展背景

云计算发展历程

  • 2006年:Amazon推出AWS EC2和S3服务,开启云计算时代
  • 2008年:Google发布Google App Engine
  • 2010年:Microsoft Azure正式上线
  • 2015年至今:云计算进入成熟期,混合云、多云成为主流

大数据技术演进

  • 2004年:Google发布MapReduce论文
  • 2006年:Apache Hadoop项目启动
  • 2009年:Apache Spark诞生
  • 2012年至今:实时计算、流处理、机器学习融合

在结构健康监测中的应用

  • 海量监测数据的存储与管理
  • 实时数据处理与异常检测
  • 长期趋势分析与性能退化评估
  • 多源数据融合与综合分析

1.2 结构健康监测中的数据挑战

数据规模挑战

  • 大型桥梁:数百至数千传感器,采样率100Hz,日数据量达TB级
  • 传感器类型:加速度、应变、位移、温度、风速等
  • 数据维度:时间序列、空间分布、多物理场

数据特征

特征描述挑战
海量性数据量巨大存储成本高,查询效率低
高速性实时数据流处理延迟要求高
多样性多源异构数据数据整合困难
价值密度低有用信息占比小数据挖掘难度大
时效性数据价值随时间递减需要实时分析

传统方案的局限

  • 本地存储容量有限
  • 计算能力不足
  • 扩展性差
  • 维护成本高
  • 数据孤岛问题

1.3 云计算与大数据的优势

云计算优势

  1. 弹性扩展:按需分配资源,应对数据峰值
  2. 成本优化:按需付费,降低基础设施投入
  3. 高可用性:多地域部署,保证服务连续性
  4. 专业运维:云服务商提供基础设施维护
  5. 全球访问:随时随地访问监测数据

大数据技术优势

  1. 分布式存储:水平扩展,支持PB级数据
  2. 并行计算:提高数据处理效率
  3. 实时处理:流式计算,毫秒级响应
  4. 智能分析:机器学习,自动模式识别
  5. 可视化展示:直观呈现分析结果

云计算基础

2.1 云计算服务模型

基础设施即服务(IaaS)

  • 提供虚拟化的计算资源
  • 用户控制操作系统和应用程序
  • 代表产品:AWS EC2、Azure VM、阿里云ECS
  • 适用场景:需要完全控制环境的应用

平台即服务(PaaS)

  • 提供应用开发和部署平台
  • 用户专注于应用开发,无需管理基础设施
  • 代表产品:AWS Elastic Beanstalk、Google App Engine
  • 适用场景:快速开发和部署监测应用

软件即服务(SaaS)

  • 提供完整的软件应用
  • 用户通过浏览器访问
  • 代表产品:监测数据管理平台、可视化系统
  • 适用场景:直接使用监测服务

函数即服务(FaaS)

  • 事件驱动的计算模型
  • 按需执行,按调用付费
  • 代表产品:AWS Lambda、Azure Functions
  • 适用场景:数据预处理、事件响应

2.2 云计算部署模式

公有云

  • 云资源由第三方提供商拥有和运营
  • 通过互联网提供服务
  • 优点:成本低、弹性好、维护简单
  • 缺点:数据安全性、合规性挑战
  • 适用:非敏感数据的存储和处理

私有云

  • 云资源仅供单一组织使用
  • 可部署在本地或托管在第三方数据中心
  • 优点:安全性高、可控性强
  • 缺点:成本高、维护复杂
  • 适用:敏感监测数据的处理

混合云

  • 公有云和私有云的组合
  • 数据和应用在两者之间流动
  • 优点:灵活性高、成本优化
  • 挑战:数据同步、安全管理
  • 适用:分级数据管理策略

多云

  • 使用多个云服务提供商
  • 避免供应商锁定
  • 优点:高可用、成本优化
  • 挑战:管理复杂度高
  • 适用:大型监测系统的容灾备份

2.3 云计算核心技术

虚拟化技术

  • 服务器虚拟化:VMware、KVM、Xen
  • 容器化:Docker、Kubernetes
  • 网络虚拟化:SDN、NFV
  • 存储虚拟化:分布式存储系统

分布式计算

  • MapReduce:批处理计算模型
  • Spark:内存计算框架
  • Flink:流处理引擎
  • Storm:实时计算系统

微服务架构

  • 服务拆分:按功能模块拆分
  • 服务治理:注册发现、负载均衡
  • 容器编排:Kubernetes、Docker Swarm
  • DevOps:CI/CD流水线

无服务器架构

  • 事件驱动:响应监测事件
  • 自动扩展:根据负载自动调整
  • 按需付费:降低空闲成本
  • 快速部署:简化运维

2.4 云原生技术

容器技术

  • Docker:应用容器化
  • 镜像管理:Docker Registry、Harbor
  • 容器运行时:containerd、CRI-O

容器编排

  • Kubernetes:容器编排标准
  • 服务网格:Istio、Linkerd
  • 自动扩缩容:HPA、VPA

持续集成/持续部署(CI/CD)

  • 代码管理:Git、GitLab
  • 构建工具:Jenkins、GitLab CI
  • 部署策略:蓝绿部署、金丝雀发布

可观测性

  • 日志收集:ELK Stack、Fluentd
  • 监控告警:Prometheus、Grafana
  • 链路追踪:Jaeger、Zipkin

大数据技术架构

3.1 大数据技术栈

数据采集层

  • 消息队列:Kafka、RabbitMQ、RocketMQ
  • 数据采集:Flume、Logstash、Filebeat
  • 数据同步:Sqoop、DataX、Canal

数据存储层

  • 分布式文件系统:HDFS、Ceph、MinIO
  • NoSQL数据库:HBase、Cassandra、MongoDB
  • 时序数据库:InfluxDB、TimescaleDB、TDengine
  • 数据仓库:Hive、ClickHouse、Doris

数据处理层

  • 批处理:MapReduce、Spark SQL
  • 流处理:Spark Streaming、Flink、Storm
  • 交互式查询:Presto、Impala、Druid

数据分析层

  • 机器学习:Spark MLlib、Flink ML
  • 深度学习:TensorFlow、PyTorch
  • 图计算:GraphX、Neo4j

数据应用层

  • 可视化:Tableau、Power BI、Grafana
  • BI工具:Superset、Metabase
  • 报表系统:自定义开发

3.2 数据采集与传输

传感器数据接入

  • MQTT协议:轻量级物联网通信
  • CoAP协议:受限环境应用协议
  • HTTP/HTTPS:RESTful API接口
  • WebSocket:实时双向通信

数据格式

  • JSON:结构化数据交换
  • Protocol Buffers:高效二进制格式
  • Avro:数据序列化系统
  • Parquet:列式存储格式

数据传输模式

  • 实时传输:数据产生后立即发送
  • 批量传输:定时批量上传
  • 断点续传:网络中断后恢复传输
  • 压缩传输:减少网络带宽占用

数据质量控制

  • 数据校验:格式、范围、完整性检查
  • 数据清洗:异常值处理、缺失值填充
  • 数据去重:消除重复数据
  • 数据标准化:统一数据格式和单位

3.3 数据存储架构

分层存储架构

热数据层(Hot)

  • 存储介质:SSD、内存
  • 数据特征:最近7天数据,高频访问
  • 应用场景:实时监控、异常检测
  • 技术选型:Redis、InfluxDB

温数据层(Warm)

  • 存储介质:SATA磁盘
  • 数据特征:7-30天数据,中频访问
  • 应用场景:趋势分析、报表生成
  • 技术选型:Elasticsearch、ClickHouse

冷数据层(Cold)

  • 存储介质:对象存储、磁带
  • 数据特征:30天以上数据,低频访问
  • 应用场景:历史查询、合规审计
  • 技术选型:HDFS、S3、OSS

数据生命周期管理

  • 自动迁移:根据访问频率自动分层
  • 数据归档:长期保存的压缩存储
  • 数据删除:过期数据的自动清理
  • 数据备份:多副本保证数据安全

3.4 实时计算架构

Lambda架构

  • 批处理层:全量数据离线处理
  • 速度层:实时数据流处理
  • 服务层:合并批处理和实时结果
  • 优点:容错性好、数据一致性
  • 缺点:维护两套系统,复杂度高

Kappa架构

  • 统一流处理:所有数据按流处理
  • 重新处理:需要时重新消费历史数据
  • 优点:架构简单、维护方便
  • 缺点:对消息队列要求高

流批一体架构

  • 统一计算引擎:Flink、Spark
  • 统一数据存储:湖仓一体
  • 统一编程模型:SQL/API
  • 优点:降低复杂度、提高开发效率

实时处理场景

  • 实时监控:结构状态实时展示
  • 异常检测:毫秒级异常识别
  • 预警通知:即时告警推送
  • 动态控制:实时反馈控制

数据存储与管理

4.1 时序数据库

时序数据特点

  • 时间戳为主键
  • 数据按时间顺序写入
  • 高写入吞吐量
  • 时间范围查询为主
  • 数据压缩率高

主流时序数据库

InfluxDB

  • 高性能写入和查询
  • 类SQL查询语言
  • 自动数据保留策略
  • 集群版商业授权

TimescaleDB

  • 基于PostgreSQL
  • 支持标准SQL
  • 自动分区管理
  • 开源免费

TDengine

  • 国产时序数据库
  • 超高性能写入
  • 边缘云协同
  • 开源+商业版

数据建模

  • 超级表:定义数据结构
  • 子表:按设备分表
  • 标签:设备属性索引
  • 数据类型:时间戳、数值、布尔、字符串

4.2 数据仓库

数据仓库架构

  • ODS层:操作数据存储,原始数据
  • DWD层:明细数据层,清洗转换
  • DWS层:汇总数据层,轻度聚合
  • ADS层:应用数据层,面向应用

Hive数据仓库

  • 基于Hadoop的数据仓库工具
  • 类SQL查询语言(HiveQL)
  • 支持多种存储格式(ORC、Parquet)
  • 适合离线批处理

ClickHouse

  • 列式存储数据库
  • 高性能OLAP查询
  • 实时数据插入
  • 适合时序数据分析

数据建模方法

  • 星型模型:事实表+维度表
  • 雪花模型:维度表规范化
  • 宽表模型:预聚合减少关联
  • 数据湖:原始数据直接存储

4.3 数据湖与湖仓一体

数据湖概念

  • 存储原始格式的海量数据
  • 支持结构化、半结构化、非结构化数据
  • Schema-on-read:读取时定义结构
  • 低成本存储:对象存储为主

数据湖技术

  • Delta Lake:ACID事务支持
  • Apache Iceberg:表格式管理
  • Apache Hudi:增量数据处理
  • 对象存储:S3、OSS、MinIO

湖仓一体(Lakehouse)

  • 数据湖的存储成本 + 数据仓库的性能
  • 统一元数据管理
  • 支持BI和ML工作负载
  • 代表产品:Databricks、Starburst

在结构健康监测中的应用

  • 原始监测数据长期保存
  • 多源异构数据整合
  • 历史数据回溯分析
  • 机器学习训练数据准备

4.4 数据治理

数据质量管理

  • 完整性:数据缺失率监控
  • 准确性:数据校验规则
  • 一致性:跨系统数据对齐
  • 时效性:数据更新延迟监控

元数据管理

  • 技术元数据:表结构、字段类型
  • 业务元数据:业务含义、数据Owner
  • 操作元数据:数据血缘、处理日志
  • 管理元数据:安全级别、保留策略

数据安全

  • 访问控制:RBAC、ABAC
  • 数据加密:传输加密、存储加密
  • 数据脱敏:敏感信息保护
  • 审计日志:操作追溯

数据标准

  • 命名规范:表名、字段名统一
  • 数据字典:标准术语定义
  • 编码规范:设备编码、位置编码
  • 质量标准:数据质量评分体系

数据处理与分析

5.1 批处理计算

MapReduce编程模型

  • Map阶段:数据分割并行处理
  • Shuffle阶段:数据排序和分组
  • Reduce阶段:结果汇总聚合
  • 适用场景:离线数据处理

Spark批处理

  • RDD:弹性分布式数据集
  • DataFrame/Dataset:结构化数据处理
  • Spark SQL:SQL查询接口
  • 内存计算:比MapReduce快10-100倍

批处理应用场景

  • 日/周/月报表生成
  • 历史数据回溯分析
  • 批量数据清洗转换
  • 模型批量训练

5.2 流处理计算

Spark Streaming

  • 微批处理模型
  • DStream:离散化流
  • Structured Streaming:结构化流处理
  • exactly-once语义

Apache Flink

  • 真正的流处理引擎
  • 事件时间处理
  • 状态管理
  • CEP复杂事件处理

流处理应用场景

  • 实时异常检测
  • 实时预警推送
  • 实时数据清洗
  • 实时聚合统计

窗口计算

  • 滚动窗口:固定大小不重叠
  • 滑动窗口:固定大小可重叠
  • 会话窗口:动态大小按活动划分
  • 全局窗口:全局统一计算

5.3 数据查询与分析

即席查询(Ad-hoc)

  • Presto/Trino:分布式SQL查询引擎
  • Impala:Cloudera的实时查询引擎
  • Drill:Schema-free SQL查询

多维分析(OLAP)

  • 上卷(Roll-up):汇总聚合
  • 下钻(Drill-down):详细展开
  • 切片(Slice):单维度筛选
  • 切块(Dice):多维度筛选
  • 旋转(Pivot):维度转换

时序数据分析

  • 降采样:降低时间分辨率
  • 插值:缺失值填充
  • 平滑:噪声去除
  • 聚合:时间窗口统计
  • 模式识别:周期性、趋势性分析

5.4 数据挖掘与机器学习

特征工程

  • 时域特征:均值、方差、峰值、有效值
  • 频域特征:频谱、功率谱密度
  • 时频特征:小波系数、STFT
  • 统计特征:偏度、峰度、熵

异常检测算法

  • 统计方法:3σ原则、箱线图
  • 机器学习方法:孤立森林、One-Class SVM
  • 深度学习方法:自编码器、LSTM
  • 集成方法:多种算法投票

预测模型

  • 时间序列预测:ARIMA、Prophet
  • 回归预测:线性回归、随机森林、XGBoost
  • 深度学习:LSTM、Transformer
  • 物理模型:有限元、数字孪生

模型部署

  • 批量预测:定时批量推理
  • 实时预测:在线推理服务
  • 模型管理:版本控制、A/B测试
  • 模型监控:性能监控、漂移检测

机器学习与人工智能

6.1 机器学习平台

开源机器学习平台

  • MLflow:机器学习生命周期管理
  • Kubeflow:Kubernetes上的机器学习
  • Airflow:工作流编排
  • Feast:特征存储

云原生机器学习

  • AWS SageMaker:全托管ML平台
  • Azure Machine Learning:微软ML云服务
  • Google AI Platform:谷歌ML平台
  • 阿里云PAI:阿里云机器学习平台

AutoML

  • 自动特征工程
  • 自动模型选择
  • 自动超参优化
  • 自动模型部署

MLOps

  • 持续集成/持续训练(CI/CT)
  • 模型版本管理
  • 实验追踪
  • 模型监控与回滚

6.2 深度学习框架

TensorFlow

  • Google开发的深度学习框架
  • 生产部署友好
  • TensorBoard可视化
  • TensorFlow Serving模型服务

PyTorch

  • Facebook开发的深度学习框架
  • 动态图机制,调试方便
  • 研究社区主流
  • TorchServe模型服务

模型优化

  • 量化:降低模型精度,减少存储
  • 剪枝:移除冗余连接
  • 蒸馏:大模型知识迁移到小模型
  • 编译优化:TensorRT、ONNX Runtime

边缘部署

  • TensorFlow Lite:移动端和嵌入式
  • ONNX:跨平台模型格式
  • OpenVINO:Intel优化工具
  • 边缘AI芯片:NVIDIA Jetson、华为昇腾

6.3 计算机视觉应用

图像分类

  • 结构表面缺陷识别
  • 损伤类型分类
  • 材料类型识别

目标检测

  • 裂缝检测
  • 锈蚀检测
  • 变形检测
  • 车辆荷载识别

语义分割

  • 损伤区域分割
  • 结构构件分割
  • 像素级精度分析

视频分析

  • 结构振动视频分析
  • 人流车流统计
  • 异常行为检测

6.4 自然语言处理应用

文本挖掘

  • 监测报告自动提取
  • 规范标准知识抽取
  • 历史事故案例分析

知识图谱

  • 结构病害知识库
  • 因果关系推理
  • 智能问答系统

文档生成

  • 自动报告生成
  • 监测日报/月报
  • 异常事件描述

可视化与交互

7.1 数据可视化技术

可视化库

  • Matplotlib:Python基础绘图库
  • Seaborn:统计可视化
  • Plotly:交互式可视化
  • ECharts:百度开源可视化库
  • D3.js:Web端数据可视化

可视化类型

  • 时序图:趋势、周期、异常
  • 散点图:相关性分析
  • 热力图:空间分布、相关性矩阵
  • 箱线图:统计分布
  • 雷达图:多维度评估
  • 桑基图:数据流向
  • 3D可视化:结构模型、模态振型

可视化设计原则

  • 简洁明了:避免过度装饰
  • 突出重点:关键信息突出
  • 交互友好:支持缩放、筛选
  • 响应式设计:适配不同设备

7.2 实时监控大屏

大屏设计

  • 布局设计:网格系统、黄金分割
  • 色彩设计:主题色、对比色、警示色
  • 动效设计:数据刷新、状态切换
  • 信息层次:核心指标、详细数据

实时数据推送

  • WebSocket:双向实时通信
  • Server-Sent Events:服务器推送
  • 轮询:定时请求更新
  • 增量更新:只传输变化数据

典型大屏场景

  • 桥梁健康状态总览
  • 结构安全预警中心
  • 监测数据统计分析
  • 设备运行状态监控

7.3 三维可视化

BIM集成

  • BIM模型导入:IFC格式解析
  • 传感器位置标注
  • 监测数据三维展示
  • 损伤位置可视化

GIS集成

  • 地图底图:高德、百度、天地图
  • 空间定位:GPS坐标映射
  • 区域热力图:监测密度展示
  • 路径规划:巡检路线优化

数字孪生

  • 三维结构模型
  • 实时数据映射
  • 物理仿真集成
  • 虚实交互控制

7.4 移动端应用

移动监测App

  • 实时数据查看
  • 告警信息推送
  • 巡检任务管理
  • 现场数据采集

跨平台开发

  • React Native:JavaScript开发
  • Flutter:Dart开发
  • 小程序:微信/支付宝
  • PWA:渐进式Web应用

离线功能

  • 数据本地缓存
  • 离线地图
  • 离线填报
  • 数据同步

安全与隐私保护

8.1 云安全架构

安全责任共担模型

  • 云服务商:基础设施安全
  • 用户:数据和应用安全
  • 共同责任:操作系统、中间件

身份与访问管理(IAM)

  • 身份认证:多因素认证(MFA)
  • 访问控制:RBAC、ABAC
  • 单点登录(SSO)
  • 临时凭证:STS、AssumeRole

网络安全

  • VPC:虚拟私有云
  • 安全组:实例级防火墙
  • 网络ACL:子网级防火墙
  • VPN:加密通信通道
  • DDoS防护:流量清洗

8.2 数据安全

数据加密

  • 传输加密:TLS/SSL
  • 存储加密:AES-256
  • 密钥管理:KMS、HSM
  • 客户端加密:端到端加密

数据脱敏

  • 静态脱敏:存储时脱敏
  • 动态脱敏:查询时脱敏
  • 格式保留加密:保持数据格式
  • 令牌化:敏感数据替换

数据备份与恢复

  • 定期备份:全量+增量
  • 跨区域复制:灾难恢复
  • 备份加密:保护备份数据
  • 恢复演练:验证备份有效性

8.3 合规与审计

数据合规

  • 等保2.0:网络安全等级保护
  • GDPR:欧盟数据保护条例
  • 数据安全法:中国数据安全法规
  • 个人信息保护法:隐私保护

安全审计

  • 操作日志:记录所有操作
  • 访问日志:记录数据访问
  • 变更审计:配置变更追踪
  • 合规报告:自动生成合规报告

安全监控

  • SIEM:安全信息和事件管理
  • 入侵检测:异常行为识别
  • 漏洞扫描:定期安全检查
  • 威胁情报:外部威胁预警

工程应用案例

9.1 大型桥梁监测云平台

系统架构

  • 边缘层:桥梁现场数据采集
  • 传输层:4G/5G/光纤网络
  • 平台层:云原生监测平台
  • 应用层:Web端和移动端应用

数据规模

  • 传感器数量:500+个
  • 数据采样率:100Hz
  • 日数据量:约500GB
  • 历史数据:PB级存储

核心功能

  • 实时监测:结构响应实时展示
  • 异常检测:AI驱动的异常识别
  • 健康评估:基于规范的安全评估
  • 预警系统:多级预警机制
  • 报告生成:自动化监测报告

技术选型

  • 云平台:阿里云/腾讯云
  • 时序数据库:TDengine
  • 流处理:Apache Flink
  • 机器学习:TensorFlow
  • 可视化:Grafana + 自定义大屏

9.2 城市级基础设施监测中心

监测范围

  • 城市桥梁:100+座
  • 隧道:50+条
  • 地铁线路:多条
  • 高层建筑:重点建筑

数据整合

  • 多源数据接入:不同厂商系统
  • 数据标准化:统一数据格式
  • 元数据管理:统一资产管理
  • 数据共享:跨部门数据交换

分析应用

  • 城市级风险评估
  • 交通荷载分析
  • 极端天气影响分析
  • 应急响应支持

9.3 风电场群监测系统

监测对象

  • 风机叶片:应变、振动、损伤
  • 塔筒:倾斜、振动
  • 基础:沉降、应力
  • 升压站:设备状态

数据特点

  • 分布广泛:数平方公里
  • 环境恶劣:高海拔、强风
  • 数据量大:高频振动数据
  • 实时性要求:故障预警

云平台功能

  • 设备健康管理
  • 发电量预测
  • 运维优化
  • 备件预测

9.4 历史建筑数字化保护

监测需求

  • 结构变形:倾斜、沉降
  • 环境因素:温湿度、光照
  • 游客影响:振动、荷载
  • 材料劣化:裂缝扩展

技术方案

  • 微型传感器:无损安装
  • 低功耗设计:长寿命运行
  • 无线传输:LoRa/NB-IoT
  • 云平台:长期数据保存

数字化应用

  • 数字档案:全生命周期记录
  • 虚拟游览:VR/AR展示
  • 保护决策:数据驱动决策
  • 公众教育:科普展示

Python仿真实现

10.1 仿真系统架构

本仿真系统实现了云计算与大数据处理的核心功能:

  1. 分布式数据存储模拟
  2. 批处理与流处理计算
  3. 数据查询与分析
  4. 机器学习模型训练与推理
  5. 可视化展示

10.2 核心代码说明

CloudStorage类

  • 模拟云存储系统
  • 支持数据分片存储
  • 数据冗余备份

BatchProcessor类

  • 批处理计算模拟
  • MapReduce模型
  • 聚合统计分析

StreamProcessor类

  • 流处理计算模拟
  • 滑动窗口计算
  • 实时异常检测

MLModel类

  • 机器学习模型
  • 训练与推理
  • 模型评估

Visualizer类

  • 数据可视化
  • 实时监控展示
  • 分析结果呈现
"""
结构健康监测仿真 - 主题044:云计算与大数据技术
Cloud Computing and Big Data for Structural Health Monitoring
"""

import numpy as np
import matplotlib.pyplot as plt
from matplotlib.patches import Rectangle, FancyBboxPatch, Circle
import matplotlib.animation as animation
from scipy import stats
from sklearn.ensemble import IsolationForest
from sklearn.preprocessing import StandardScaler
import warnings
warnings.filterwarnings('ignore')

# 设置中文字体
plt.rcParams['font.sans-serif'] = ['SimHei', 'DejaVu Sans']
plt.rcParams['axes.unicode_minus'] = False

# 使用Agg后端,不弹出窗口
plt.switch_backend('Agg')


class CloudStorage:
    """云存储系统模拟"""
    
    def __init__(self, n_nodes=5, replication_factor=3):
        """
        初始化云存储系统
        
        Parameters:
        -----------
        n_nodes : int
            存储节点数量
        replication_factor : int
            数据副本数
        """
        self.n_nodes = n_nodes
        self.replication_factor = replication_factor
        self.nodes = []
        self.total_capacity = 0
        self.used_capacity = 0
        
        # 初始化存储节点
        for i in range(n_nodes):
            node = {
                'id': i,
                'capacity': 1000,  # GB
                'used': 0,
                'status': 'active'
            }
            self.nodes.append(node)
            self.total_capacity += node['capacity']
    
    def store_data(self, data_size, data_id):
        """存储数据"""
        # 选择存储节点(简单哈希)
        primary_node = data_id % self.n_nodes
        
        # 复制到多个节点
        stored_nodes = []
        for i in range(self.replication_factor):
            node_id = (primary_node + i) % self.n_nodes
            if self.nodes[node_id]['used'] + data_size <= self.nodes[node_id]['capacity']:
                self.nodes[node_id]['used'] += data_size
                stored_nodes.append(node_id)
        
        self.used_capacity += data_size
        return stored_nodes
    
    def get_storage_stats(self):
        """获取存储统计"""
        return {
            'total_capacity': self.total_capacity,
            'used_capacity': self.used_capacity,
            'available_capacity': self.total_capacity - self.used_capacity,
            'utilization': self.used_capacity / self.total_capacity * 100,
            'node_stats': [(n['used'], n['capacity']) for n in self.nodes]
        }


class BatchProcessor:
    """批处理计算模拟"""
    
    def __init__(self, n_mappers=4, n_reducers=2):
        """
        初始化批处理器
        
        Parameters:
        -----------
        n_mappers : int
            Mapper数量
        n_reducers : int
            Reducer数量
        """
        self.n_mappers = n_mappers
        self.n_reducers = n_reducers
        
    def map_reduce(self, data, map_func, reduce_func):
        """
        模拟MapReduce计算
        
        Parameters:
        -----------
        data : list
            输入数据
        map_func : function
            Map函数
        reduce_func : function
            Reduce函数
        """
        # Map阶段
        mapped_data = []
        chunk_size = len(data) // self.n_mappers
        
        for i in range(self.n_mappers):
            start = i * chunk_size
            end = start + chunk_size if i < self.n_mappers - 1 else len(data)
            chunk = data[start:end]
            mapped_data.extend(map_func(chunk))
        
        # Shuffle阶段(按键分组)
        shuffled = {}
        for key, value in mapped_data:
            if key not in shuffled:
                shuffled[key] = []
            shuffled[key].append(value)
        
        # Reduce阶段
        results = []
        for key, values in shuffled.items():
            result = reduce_func(key, values)
            results.append((key, result))
        
        return results
    
    def aggregate_statistics(self, sensor_data):
        """聚合统计计算"""
        # Map: 提取统计信息
        def map_func(chunk):
            results = []
            for record in chunk:
                sensor_id = record['sensor_id']
                value = record['value']
                results.append((sensor_id, value))
            return results
        
        # Reduce: 计算统计量
        def reduce_func(sensor_id, values):
            return {
                'count': len(values),
                'mean': np.mean(values),
                'std': np.std(values),
                'min': np.min(values),
                'max': np.max(values)
            }
        
        return self.map_reduce(sensor_data, map_func, reduce_func)


class StreamProcessor:
    """流处理计算模拟"""
    
    def __init__(self, window_size=100, slide_size=10):
        """
        初始化流处理器
        
        Parameters:
        -----------
        window_size : int
            窗口大小
        slide_size : int
            滑动步长
        """
        self.window_size = window_size
        self.slide_size = slide_size
        self.buffer = []
        self.anomalies = []
        
    def process_stream(self, data_stream):
        """处理数据流"""
        results = []
        
        for i, data_point in enumerate(data_stream):
            self.buffer.append(data_point)
            
            # 当缓冲区达到窗口大小时进行处理
            if len(self.buffer) >= self.window_size:
                window_data = self.buffer[-self.window_size:]
                
                # 计算窗口统计
                window_stats = {
                    'timestamp': i,
                    'mean': np.mean(window_data),
                    'std': np.std(window_data),
                    'max': np.max(window_data),
                    'min': np.min(window_data)
                }
                
                # 异常检测(3σ原则)
                mean = window_stats['mean']
                std = window_stats['std']
                if std > 0:
                    z_scores = [(x - mean) / std for x in window_data]
                    if any(abs(z) > 3 for z in z_scores):
                        window_stats['anomaly'] = True
                        self.anomalies.append(i)
                    else:
                        window_stats['anomaly'] = False
                
                results.append(window_stats)
                
                # 滑动窗口
                if len(self.buffer) > self.window_size + self.slide_size:
                    self.buffer = self.buffer[self.slide_size:]
        
        return results


class MLModel:
    """机器学习模型"""
    
    def __init__(self, model_type='isolation_forest'):
        """
        初始化机器学习模型
        
        Parameters:
        -----------
        model_type : str
            模型类型
        """
        self.model_type = model_type
        self.model = None
        self.scaler = StandardScaler()
        
    def train(self, X_train):
        """训练模型"""
        X_scaled = self.scaler.fit_transform(X_train)
        
        if self.model_type == 'isolation_forest':
            self.model = IsolationForest(contamination=0.1, random_state=42)
            self.model.fit(X_scaled)
        
        return self
    
    def predict(self, X):
        """预测"""
        if self.model is None:
            raise ValueError("Model not trained")
        
        X_scaled = self.scaler.transform(X)
        predictions = self.model.predict(X_scaled)
        scores = self.model.score_samples(X_scaled)
        
        return predictions, scores
    
    def evaluate(self, X_test, y_true):
        """评估模型"""
        predictions, _ = self.predict(X_test)
        
        # 计算准确率
        accuracy = np.mean(predictions == y_true)
        
        # 计算精确率、召回率
        tp = np.sum((predictions == -1) & (y_true == -1))
        fp = np.sum((predictions == -1) & (y_true == 1))
        fn = np.sum((predictions == 1) & (y_true == -1))
        
        precision = tp / (tp + fp) if (tp + fp) > 0 else 0
        recall = tp / (tp + fn) if (tp + fn) > 0 else 0
        f1 = 2 * precision * recall / (precision + recall) if (precision + recall) > 0 else 0
        
        return {
            'accuracy': accuracy,
            'precision': precision,
            'recall': recall,
            'f1_score': f1
        }


class DataLake:
    """数据湖模拟"""
    
    def __init__(self):
        """初始化数据湖"""
        self.raw_data = []
        self.processed_data = []
        self.metadata = {}
        
    def ingest(self, data, source, timestamp):
        """数据接入"""
        record = {
            'data': data,
            'source': source,
            'ingest_time': timestamp,
            'format': 'raw'
        }
        self.raw_data.append(record)
        
    def process(self, processing_func):
        """数据处理"""
        for record in self.raw_data:
            processed = processing_func(record['data'])
            self.processed_data.append({
                'original': record,
                'processed': processed,
                'process_time': len(self.processed_data)
            })
        
    def query(self, condition):
        """数据查询"""
        results = []
        for record in self.processed_data:
            if condition(record):
                results.append(record)
        return results


def generate_sensor_data(n_samples=10000, n_sensors=50):
    """生成模拟传感器数据"""
    data = []
    
    for i in range(n_samples):
        timestamp = i * 0.01  # 100Hz采样
        sensor_id = i % n_sensors
        
        # 基础信号 + 噪声 + 偶尔异常
        base_signal = np.sin(2 * np.pi * 2 * timestamp)  # 2Hz正弦
        noise = np.random.normal(0, 0.1)
        
        # 10%概率出现异常
        if np.random.random() < 0.1:
            anomaly = np.random.choice([-1, 1]) * np.random.uniform(2, 5)
        else:
            anomaly = 0
        
        value = base_signal + noise + anomaly
        
        data.append({
            'timestamp': timestamp,
            'sensor_id': sensor_id,
            'value': value,
            'is_anomaly': anomaly != 0
        })
    
    return data


def simulate_cloud_bigdata():
    """云计算与大数据仿真"""
    print("="*60)
    print("云计算与大数据技术仿真 - 结构健康监测")
    print("="*60)
    
    # 生成模拟数据
    print("\n1. 生成模拟传感器数据...")
    sensor_data = generate_sensor_data(n_samples=10000, n_sensors=50)
    print(f"   生成了 {len(sensor_data)} 条传感器记录")
    
    # 云存储模拟
    print("\n2. 云存储系统模拟...")
    storage = CloudStorage(n_nodes=5, replication_factor=3)
    
    # 存储数据
    for i, record in enumerate(sensor_data):
        data_size = 0.001  # 每条记录约1KB
        storage.store_data(data_size, i)
    
    storage_stats = storage.get_storage_stats()
    print(f"   总容量: {storage_stats['total_capacity']} GB")
    print(f"   已使用: {storage_stats['used_capacity']:.2f} GB")
    print(f"   利用率: {storage_stats['utilization']:.2f}%")
    
    # 批处理计算
    print("\n3. 批处理计算(MapReduce)...")
    batch_processor = BatchProcessor(n_mappers=4, n_reducers=2)
    stats_results = batch_processor.aggregate_statistics(sensor_data)
    print(f"   处理了 {len(stats_results)} 个传感器的统计信息")
    
    # 流处理计算
    print("\n4. 流处理计算...")
    stream_processor = StreamProcessor(window_size=100, slide_size=10)
    values = [d['value'] for d in sensor_data]
    stream_results = stream_processor.process_stream(values)
    n_anomalies_detected = len(stream_processor.anomalies)
    print(f"   处理了 {len(stream_results)} 个窗口")
    print(f"   检测到 {n_anomalies_detected} 个异常")
    
    # 机器学习异常检测
    print("\n5. 机器学习异常检测...")
    # 准备特征
    features = []
    labels = []
    
    for i in range(0, len(sensor_data) - 10, 10):
        window = sensor_data[i:i+10]
        feature = [
            np.mean([d['value'] for d in window]),
            np.std([d['value'] for d in window]),
            np.max([d['value'] for d in window]),
            np.min([d['value'] for d in window]),
            np.mean([d['value']**2 for d in window])
        ]
        features.append(feature)
        # 如果窗口中有异常,标记为异常
        is_anomaly = any(d['is_anomaly'] for d in window)
        labels.append(-1 if is_anomaly else 1)
    
    features = np.array(features)
    labels = np.array(labels)
    
    # 划分训练集和测试集
    split_idx = int(len(features) * 0.8)
    X_train, X_test = features[:split_idx], features[split_idx:]
    y_train, y_test = labels[:split_idx], labels[split_idx:]
    
    # 训练模型
    ml_model = MLModel(model_type='isolation_forest')
    ml_model.train(X_train)
    
    # 评估模型
    metrics = ml_model.evaluate(X_test, y_test)
    print(f"   准确率: {metrics['accuracy']:.3f}")
    print(f"   精确率: {metrics['precision']:.3f}")
    print(f"   召回率: {metrics['recall']:.3f}")
    print(f"   F1分数: {metrics['f1_score']:.3f}")
    
    # 数据湖模拟
    print("\n6. 数据湖模拟...")
    data_lake = DataLake()
    
    # 接入数据
    for record in sensor_data[:1000]:
        data_lake.ingest(record['value'], f"sensor_{record['sensor_id']}", record['timestamp'])
    
    print(f"   接入了 {len(data_lake.raw_data)} 条原始数据")
    
    return {
        'sensor_data': sensor_data,
        'storage_stats': storage_stats,
        'batch_results': stats_results,
        'stream_results': stream_results,
        'n_anomalies_detected': n_anomalies_detected,
        'ml_metrics': metrics,
        'ml_model': ml_model,
        'features': features,
        'labels': labels
    }


def create_visualization(results):
    """创建可视化"""
    print("\n生成可视化...")
    
    fig = plt.figure(figsize=(20, 16))
    
    sensor_data = results['sensor_data']
    storage_stats = results['storage_stats']
    batch_results = results['batch_results']
    stream_results = results['stream_results']
    features = results['features']
    labels = results['labels']
    
    # 1. 传感器数据时序图
    ax1 = plt.subplot(4, 4, 1)
    timestamps = [d['timestamp'] for d in sensor_data[:1000]]
    values = [d['value'] for d in sensor_data[:1000]]
    anomalies = [i for i, d in enumerate(sensor_data[:1000]) if d['is_anomaly']]
    
    ax1.plot(timestamps, values, 'b-', linewidth=0.5, alpha=0.7)
    if anomalies:
        ax1.scatter([timestamps[i] for i in anomalies], 
                   [values[i] for i in anomalies],
                   c='red', s=20, zorder=5, label='Anomaly')
    ax1.set_xlabel('Time (s)')
    ax1.set_ylabel('Sensor Value')
    ax1.set_title('Sensor Data Time Series')
    ax1.legend()
    ax1.grid(True, alpha=0.3)
    
    # 2. 云存储节点分布
    ax2 = plt.subplot(4, 4, 2)
    node_ids = range(len(storage_stats['node_stats']))
    used = [s[0] for s in storage_stats['node_stats']]
    capacity = [s[1] for s in storage_stats['node_stats']]
    
    x = np.arange(len(node_ids))
    width = 0.35
    ax2.bar(x - width/2, capacity, width, label='Capacity', alpha=0.7, color='lightblue')
    ax2.bar(x + width/2, used, width, label='Used', alpha=0.7, color='orange')
    ax2.set_xlabel('Node ID')
    ax2.set_ylabel('Storage (GB)')
    ax2.set_title('Cloud Storage Distribution')
    ax2.set_xticks(x)
    ax2.set_xticklabels([f'Node {i}' for i in node_ids])
    ax2.legend()
    ax2.grid(True, alpha=0.3)
    
    # 3. 存储利用率
    ax3 = plt.subplot(4, 4, 3)
    utilization = storage_stats['utilization']
    colors = ['green' if utilization < 50 else 'yellow' if utilization < 80 else 'red']
    ax3.bar(['Storage Utilization'], [utilization], color=colors, alpha=0.7)
    ax3.axhline(y=50, color='orange', linestyle='--', label='50% threshold')
    ax3.axhline(y=80, color='red', linestyle='--', label='80% threshold')
    ax3.set_ylabel('Utilization (%)')
    ax3.set_title('Storage Utilization')
    ax3.legend()
    ax3.set_ylim(0, 100)
    ax3.grid(True, alpha=0.3)
    
    # 4. 批处理统计结果
    ax4 = plt.subplot(4, 4, 4)
    sensor_ids = [r[0] for r in batch_results[:20]]
    means = [r[1]['mean'] for r in batch_results[:20]]
    stds = [r[1]['std'] for r in batch_results[:20]]
    
    ax4.errorbar(sensor_ids, means, yerr=stds, fmt='o', capsize=5, alpha=0.7)
    ax4.set_xlabel('Sensor ID')
    ax4.set_ylabel('Mean ± Std')
    ax4.set_title('Batch Processing: Sensor Statistics')
    ax4.grid(True, alpha=0.3)
    
    # 5. 流处理窗口统计
    ax5 = plt.subplot(4, 4, 5)
    window_timestamps = [r['timestamp'] for r in stream_results]
    window_means = [r['mean'] for r in stream_results]
    window_stds = [r['std'] for r in stream_results]
    
    ax5.plot(window_timestamps, window_means, 'b-', label='Mean', alpha=0.7)
    ax5.fill_between(window_timestamps, 
                     np.array(window_means) - np.array(window_stds),
                     np.array(window_means) + np.array(window_stds),
                     alpha=0.3, label='±1 Std')
    ax5.set_xlabel('Window Index')
    ax5.set_ylabel('Value')
    ax5.set_title('Stream Processing: Window Statistics')
    ax5.legend()
    ax5.grid(True, alpha=0.3)
    
    # 6. 异常检测结果
    ax6 = plt.subplot(4, 4, 6)
    anomaly_flags = [r['anomaly'] for r in stream_results]
    anomaly_indices = [i for i, flag in enumerate(anomaly_flags) if flag]
    
    ax6.plot(window_timestamps, window_means, 'b-', alpha=0.5)
    if anomaly_indices:
        ax6.scatter([window_timestamps[i] for i in anomaly_indices],
                   [window_means[i] for i in anomaly_indices],
                   c='red', s=50, zorder=5, label='Detected Anomaly')
    ax6.set_xlabel('Window Index')
    ax6.set_ylabel('Mean Value')
    ax6.set_title('Anomaly Detection Results')
    ax6.legend()
    ax6.grid(True, alpha=0.3)
    
    # 7. 特征分布
    ax7 = plt.subplot(4, 4, 7)
    feature_names = ['Mean', 'Std', 'Max', 'Min', 'RMS']
    normal_features = features[labels == 1]
    anomaly_features = features[labels == -1]
    
    x_pos = np.arange(len(feature_names))
    
    # 正常数据箱线图
    bp1 = ax7.boxplot([normal_features[:, i] for i in range(5)], 
               positions=x_pos - 0.2, widths=0.3, patch_artist=True,
               boxprops=dict(facecolor='lightblue', alpha=0.7))
    
    # 异常数据箱线图
    if len(anomaly_features) > 0:
        bp2 = ax7.boxplot([anomaly_features[:, i] for i in range(5)],
                   positions=x_pos + 0.2, widths=0.3, patch_artist=True,
                   boxprops=dict(facecolor='lightcoral', alpha=0.7))
    
    ax7.set_xticks(x_pos)
    ax7.set_xticklabels(feature_names)
    ax7.set_ylabel('Feature Value')
    ax7.set_title('Feature Distribution')
    ax7.grid(True, alpha=0.3)
    
    # 添加图例
    from matplotlib.patches import Patch
    legend_elements = [Patch(facecolor='lightblue', alpha=0.7, label='Normal'),
                      Patch(facecolor='lightcoral', alpha=0.7, label='Anomaly')]
    ax7.legend(handles=legend_elements, loc='upper right')
    
    # 8. 机器学习性能指标
    ax8 = plt.subplot(4, 4, 8)
    metrics = results['ml_metrics']
    metric_names = ['Accuracy', 'Precision', 'Recall', 'F1-Score']
    metric_values = [metrics['accuracy'], metrics['precision'], 
                    metrics['recall'], metrics['f1_score']]
    
    bars = ax8.bar(metric_names, metric_values, color=['green', 'blue', 'orange', 'red'], alpha=0.7)
    ax8.set_ylabel('Score')
    ax8.set_title('ML Model Performance')
    ax8.set_ylim(0, 1)
    ax8.grid(True, alpha=0.3)
    
    # 添加数值标签
    for bar, value in zip(bars, metric_values):
        height = bar.get_height()
        ax8.text(bar.get_x() + bar.get_width()/2., height,
                f'{value:.3f}', ha='center', va='bottom')
    
    # 9. 数据分布直方图
    ax9 = plt.subplot(4, 4, 9)
    all_values = [d['value'] for d in sensor_data]
    ax9.hist(all_values, bins=50, color='skyblue', edgecolor='black', alpha=0.7)
    ax9.axvline(np.mean(all_values), color='red', linestyle='--', 
               label=f'Mean: {np.mean(all_values):.3f}')
    ax9.axvline(np.median(all_values), color='green', linestyle='--',
               label=f'Median: {np.median(all_values):.3f}')
    ax9.set_xlabel('Value')
    ax9.set_ylabel('Frequency')
    ax9.set_title('Data Distribution')
    ax9.legend()
    ax9.grid(True, alpha=0.3)
    
    # 10. 传感器数据热力图
    ax10 = plt.subplot(4, 4, 10)
    # 重塑数据为传感器 x 时间格式
    sensor_matrix = np.array(values[:500]).reshape(10, 50)
    im = ax10.imshow(sensor_matrix, aspect='auto', cmap='viridis')
    ax10.set_xlabel('Time Index')
    ax10.set_ylabel('Sensor Group')
    ax10.set_title('Sensor Data Heatmap')
    plt.colorbar(im, ax=ax10, label='Value')
    
    # 11. 处理延迟分析
    ax11 = plt.subplot(4, 4, 11)
    # 模拟处理延迟
    batch_delays = np.random.exponential(2, 100)  # 批处理延迟
    stream_delays = np.random.exponential(0.1, 100)  # 流处理延迟
    
    ax11.hist(batch_delays, bins=20, alpha=0.5, label='Batch Processing', color='blue')
    ax11.hist(stream_delays, bins=20, alpha=0.5, label='Stream Processing', color='orange')
    ax11.set_xlabel('Processing Delay (s)')
    ax11.set_ylabel('Frequency')
    ax11.set_title('Processing Latency Distribution')
    ax11.legend()
    ax11.grid(True, alpha=0.3)
    
    # 12. 数据增长率
    ax12 = plt.subplot(4, 4, 12)
    time_points = np.arange(24)  # 24小时
    data_growth = np.cumsum(np.random.exponential(10, 24))  # 每小时数据增长
    
    ax12.plot(time_points, data_growth, 'g-', linewidth=2)
    ax12.fill_between(time_points, data_growth, alpha=0.3, color='green')
    ax12.set_xlabel('Hour')
    ax12.set_ylabel('Cumulative Data (GB)')
    ax12.set_title('Data Growth Over Time')
    ax12.grid(True, alpha=0.3)
    
    # 13. 资源利用率趋势
    ax13 = plt.subplot(4, 4, 13)
    hours = np.arange(24)
    cpu_usage = 30 + 40 * np.sin(2 * np.pi * hours / 24) + np.random.normal(0, 5, 24)
    memory_usage = 40 + 30 * np.sin(2 * np.pi * (hours - 6) / 24) + np.random.normal(0, 5, 24)
    
    ax13.plot(hours, cpu_usage, 'b-', label='CPU Usage', linewidth=2)
    ax13.plot(hours, memory_usage, 'r-', label='Memory Usage', linewidth=2)
    ax13.fill_between(hours, cpu_usage, alpha=0.2, color='blue')
    ax13.fill_between(hours, memory_usage, alpha=0.2, color='red')
    ax13.set_xlabel('Hour of Day')
    ax13.set_ylabel('Usage (%)')
    ax13.set_title('Resource Utilization Trend')
    ax13.legend()
    ax13.grid(True, alpha=0.3)
    
    # 14. 数据质量指标
    ax14 = plt.subplot(4, 4, 14)
    quality_metrics = ['Completeness', 'Accuracy', 'Consistency', 'Timeliness']
    quality_scores = [95, 92, 88, 90]
    colors = ['green' if s >= 90 else 'yellow' if s >= 80 else 'red' for s in quality_scores]
    
    bars = ax14.barh(quality_metrics, quality_scores, color=colors, alpha=0.7)
    ax14.set_xlabel('Score (%)')
    ax14.set_title('Data Quality Metrics')
    ax14.set_xlim(0, 100)
    ax14.grid(True, alpha=0.3, axis='x')
    
    # 15. 成本分析
    ax15 = plt.subplot(4, 4, 15)
    cost_categories = ['Storage', 'Compute', 'Network', 'ML']
    costs = [35, 40, 15, 10]
    colors = ['#ff9999', '#66b3ff', '#99ff99', '#ffcc99']
    
    wedges, texts, autotexts = ax15.pie(costs, labels=cost_categories, autopct='%1.1f%%',
                                        colors=colors, startangle=90)
    ax15.set_title('Cloud Cost Distribution')
    
    # 16. 系统统计信息
    ax16 = plt.subplot(4, 4, 16)
    ax16.axis('off')
    
    stats_text = f"""
    Cloud & Big Data Simulation
    ============================
    
    Data Statistics:
    • Total Records: {len(sensor_data):,}
    • Sensors: 50
    • Sampling Rate: 100 Hz
    • Anomaly Rate: 10%
    
    Storage Statistics:
    • Total Capacity: {storage_stats['total_capacity']} GB
    • Used Capacity: {storage_stats['used_capacity']:.2f} GB
    • Available: {storage_stats['available_capacity']:.2f} GB
    • Utilization: {storage_stats['utilization']:.2f}%
    • Replication Factor: 3
    
    Processing Statistics:
    • Batch Jobs: {len(batch_results)}
    • Stream Windows: {len(stream_results)}
    • Anomalies Detected: {results['n_anomalies_detected']}
    
    ML Performance:
    • Accuracy: {results['ml_metrics']['accuracy']:.3f}
    • Precision: {results['ml_metrics']['precision']:.3f}
    • Recall: {results['ml_metrics']['recall']:.3f}
    • F1-Score: {results['ml_metrics']['f1_score']:.3f}
    
    System Architecture:
    • Storage Nodes: 5
    • Mappers: 4
    • Reducers: 2
    • Window Size: 100
    • Slide Size: 10
    """
    
    ax16.text(0.1, 0.5, stats_text, fontsize=9, verticalalignment='center',
             fontfamily='monospace', bbox=dict(boxstyle='round', facecolor='wheat', alpha=0.5))
    
    plt.suptitle('Cloud Computing and Big Data for Structural Health Monitoring', 
                fontsize=16, fontweight='bold', y=0.995)
    plt.tight_layout(rect=[0, 0, 1, 0.99])
    plt.savefig('cloud_bigdata_analysis.png', dpi=150, bbox_inches='tight')
    print("  综合分析图已保存: cloud_bigdata_analysis.png")
    plt.close()


def create_animation(results):
    """创建数据流动画"""
    print("\n生成动画...")
    
    fig, axes = plt.subplots(1, 2, figsize=(16, 8))
    
    sensor_data = results['sensor_data']
    
    # 左图:实时数据流
    ax1 = axes[0]
    ax1.set_xlim(0, 10)
    ax1.set_ylim(-3, 3)
    ax1.set_xlabel('Time (s)', fontsize=12)
    ax1.set_ylabel('Sensor Value', fontsize=12)
    ax1.set_title('Real-time Data Stream', fontsize=14, fontweight='bold')
    ax1.grid(True, alpha=0.3)
    
    # 初始化线条
    line, = ax1.plot([], [], 'b-', linewidth=1.5, alpha=0.7)
    scatter = ax1.scatter([], [], c='red', s=30, zorder=5)
    
    # 右图:处理统计
    ax2 = axes[1]
    ax2.set_xlim(0, 100)
    ax2.set_ylim(0, 100)
    ax2.set_xlabel('Time Index', fontsize=12)
    ax2.set_ylabel('Count', fontsize=12)
    ax2.set_title('Processing Statistics', fontsize=14, fontweight='bold')
    ax2.grid(True, alpha=0.3)
    
    line_processed, = ax2.plot([], [], 'g-', linewidth=2, label='Processed')
    line_anomaly, = ax2.plot([], [], 'r-', linewidth=2, label='Anomalies')
    ax2.legend(loc='upper left')
    
    # 添加文本信息
    text_info = ax2.text(0.02, 0.98, '', transform=ax2.transAxes, fontsize=10,
                        verticalalignment='top', fontfamily='monospace',
                        bbox=dict(boxstyle='round', facecolor='wheat', alpha=0.8))
    
    # 数据准备
    window_size = 200
    n_frames = 100
    
    def init():
        line.set_data([], [])
        scatter.set_offsets(np.empty((0, 2)))
        line_processed.set_data([], [])
        line_anomaly.set_data([], [])
        text_info.set_text('Initializing...')
        return line, scatter, line_processed, line_anomaly, text_info
    
    def update(frame):
        start_idx = frame * 10
        end_idx = start_idx + window_size
        
        if end_idx > len(sensor_data):
            end_idx = len(sensor_data)
        
        # 更新左图
        window_data = sensor_data[start_idx:end_idx]
        timestamps = [d['timestamp'] for d in window_data]
        values = [d['value'] for d in window_data]
        
        line.set_data(timestamps, values)
        
        # 更新异常点
        anomaly_indices = [i for i, d in enumerate(window_data) if d['is_anomaly']]
        if anomaly_indices:
            anomaly_times = [timestamps[i] for i in anomaly_indices]
            anomaly_values = [values[i] for i in anomaly_indices]
            scatter.set_offsets(np.column_stack([anomaly_times, anomaly_values]))
        else:
            scatter.set_offsets(np.empty((0, 2)))
        
        # 更新右图
        processed_count = list(range(frame + 1))
        anomaly_count = []
        
        count = 0
        for i in range(frame + 1):
            idx = i * 10
            if idx < len(sensor_data):
                window = sensor_data[idx:idx+window_size]
                n_anomalies = sum(1 for d in window if d['is_anomaly'])
                count += n_anomalies
            anomaly_count.append(count)
        
        line_processed.set_data(processed_count, [p * 10 for p in processed_count])
        line_anomaly.set_data(processed_count, anomaly_count)
        
        # 更新文本
        current_anomalies = sum(1 for d in window_data if d['is_anomaly'])
        info_text = f"""Frame: {frame}
Data Points: {len(window_data)}
Anomalies: {current_anomalies}
Total Processed: {(frame + 1) * 10}"""
        text_info.set_text(info_text)
        
        return line, scatter, line_processed, line_anomaly, text_info
    
    # 创建动画
    anim = animation.FuncAnimation(fig, update, init_func=init,
                                  frames=n_frames, interval=200, blit=True)
    
    # 保存动画
    anim.save('cloud_bigdata_animation.gif', writer='pillow', fps=5, dpi=100)
    print("  动画已保存: cloud_bigdata_animation.gif")
    plt.close()


def main():
    """主函数"""
    print("\n" + "="*60)
    print("结构健康监测中的云计算与大数据技术仿真")
    print("="*60)
    
    # 运行仿真
    results = simulate_cloud_bigdata()
    
    # 创建可视化
    create_visualization(results)
    
    # 创建动画
    create_animation(results)
    
    print("\n" + "="*60)
    print("仿真完成!")
    print("="*60)
    print("\n生成的文件:")
    print("  - cloud_bigdata_analysis.png (综合分析图)")
    print("  - cloud_bigdata_animation.gif (数据流动画)")


if __name__ == "__main__":
    main()

10.3 运行仿真程序

python run_simulation.py

程序将生成:

  1. 综合分析图(16个子图)
  2. 动画演示(GIF格式)

10.4 仿真结果分析

通过仿真可以观察到:

  • 数据存储分布与冗余
  • 批处理与流处理性能
  • 查询响应时间
  • 机器学习模型效果

更多推荐