汽车大数据分析系统#
项目概述#
汽车大数据分析系统是一个综合性的企业级大数据解决方案,旨在通过收集、存储和分析汽车行业的各类数据,为汽车制造商、经销商和消费者提供有价值的洞察。本系统整合了Spark、Hadoop、Hive、HBase、Kafka等大数据技术,实现了从数据采集到分析、可视化的完整流程。
项目采用Lambda架构实现批处理与流处理的协同,通过Spark Streaming微批次处理实现实时数据分析,结合HBase存储海量传感器时序数据、MySQL存储业务结构化数据,构建了完整的车联网数据处理平台。
知识点讲解#
知识点1:车联网数据特征#
车联网(Internet of Vehicles, IoV)是物联网在交通领域的典型应用,其数据具有鲜明的行业特征。理解这些特征是设计合理数据处理架构的前提。
数据类型#
| 数据类型 | 来源 | 数据格式 | 采集频率 | 典型用途 |
|---|---|---|---|---|
| 传感器数据 | 发动机温度、机油压力、冷却液温度、轮胎压力、燃油液位、电池电压 | 结构化数值 | 1-100Hz | 实时监控、故障预警 |
| GPS轨迹数据 | 车载GPS模块 | 经纬度+时间戳 | 1-10Hz | 路径规划、行为分析 |
| OBD诊断数据 | 车载OBD-II接口 | DTC故障码+参数 | 事件触发 | 故障诊断、维修预测 |
| 驾驶行为数据 | 加速度计、陀螺仪、方向盘传感器 | 多维向量 | 10-100Hz | 驾驶评分、保险定价 |
数据特征#
| 特征维度 | 描述 | 技术挑战 | 应对策略 |
|---|---|---|---|
| 高吞吐 | 单车每日可产生数GB数据,百万级车队每日PB级 | 存储、传输、处理压力巨大 | 分层存储、数据压缩、边缘计算 |
| 实时性 | 安全相关数据需毫秒级响应(如碰撞预警) | 低延迟处理要求 | 流式计算、内存计算 |
| 多源异构 | 不同品牌、不同车型数据格式不统一 | 数据整合困难 | 统一数据模型、Schema Registry |
| 时序性 | 数据带有严格的时间顺序 | 时序查询、窗口计算 | 时序数据库、窗口函数 |
| 价值密度低 | 大量正常数据中蕴含少量异常信息 | 有效信息提取 | 异常检测算法、采样策略 |
类比理解:车联网数据就像一个大型医院的监护系统——每张病床(每辆车)上的各种监护仪(传感器)持续产生心率、血压、血氧等数据(传感器数据),护士站(数据处理中心)需要实时监控所有病床的状态,一旦某项指标异常(故障预警),必须立即响应。而每天的例行查房记录(批处理分析)则用于长期健康趋势分析。
知识点2:Lambda架构与Kappa架构#
在车联网数据处理中,架构选型是系统设计的核心决策。当前业界主要有两种大数据处理架构:Lambda架构和Kappa架构。
Lambda架构#
+------------------------------------------------------------------+
| Lambda 架构 |
+------------------------------------------------------------------+
| |
| +----------------+ +------------------+ +--------------+ |
| | 批处理层 | | 速度层 | | 服务层 | |
| | (Batch Layer) | | (Speed Layer) | | (Serving | |
| | | | | | Layer) | |
| | - HDFS原始数据 | | - Spark Streaming| | - 查询合并 | |
| | - Spark批处理 | | - 实时计算 | | - MySQL | |
| | - 全量重算 | | - 增量计算 | | - HBase | |
| +-------+--------+ +--------+---------+ +------+-------+ |
| | | | |
| +-------+------+--------+ | |
| | | | |
| +-------v------v--------+ | |
| | 数据输入层 +-----------------------+ |
| | (Kafka消息队列) | |
| +----------------------+ |
+------------------------------------------------------------------+- 批处理层(Batch Layer):存储全量原始数据,定期(如每天)进行全量计算,结果准确但延迟高
- 速度层(Speed Layer):实时处理新到达的数据,增量计算,延迟低但结果可能不精确
- 服务层(Serving Layer):合并批处理层和速度层的结果,对外提供统一查询接口
Kappa架构#
+------------------------------------------------------------------+
| Kappa 架构 |
+------------------------------------------------------------------+
| |
| +------------------+ +------------------+ |
| | 统一处理层 | | 服务层 | |
| | (Stream Layer) | | (Serving Layer) | |
| | | | | |
| | - Kafka保留日志 | | - 实时查询 | |
| | - Spark Streaming| | - MySQL | |
| | - 流批一体 | | - HBase | |
| +--------+---------+ +--------+---------+ |
| | | |
| +----------+-------------+ |
| | |
| +----------v------------+ |
| | 数据输入层 | |
| | (Kafka消息队列) | |
| +----------------------+ |
+------------------------------------------------------------------+- 统一处理层:所有数据处理都通过流式计算完成,通过回放Kafka历史数据实现"批处理"
- 服务层:直接查询流处理结果
架构对比#
| 对比维度 | Lambda架构 | Kappa架构 |
|---|---|---|
| 核心思想 | 批流分离,各取所长 | 流批一体,统一处理 |
| 代码复杂度 | 高(需维护批处理和流处理两套代码) | 低(只需维护一套流处理代码) |
| 数据一致性 | 需要合并批处理和速度层结果 | 天然一致(单一数据源) |
| 计算延迟 | 批处理延迟高(小时级),速度层延迟低(秒级) | 统一低延迟(秒级) |
| 历史重算 | 批处理层自动全量重算 | 通过Kafka日志回放重算 |
| 运维复杂度 | 高(两套系统运维) | 低(一套系统运维) |
| 容错能力 | 批处理层可修正速度层错误 | 依赖Kafka消息保留策略 |
| 适用场景 | 数据准确性要求极高、历史数据频繁重算 | 实时性要求高、业务逻辑相对简单 |
| 技术门槛 | 中等 | 较高(需深入理解流处理) |
本项目架构选择#
本项目采用Lambda架构,原因如下:
- 汽车销售数据分析需要精确的历史统计(批处理层保障准确性)
- 传感器监控需要实时告警(速度层保障实时性)
- 两种数据处理逻辑差异较大,分离实现更清晰
- 便于初学者理解批处理与流处理的区别
知识点3:Spark Streaming微批次处理原理#
Spark Streaming是Spark生态中的流处理组件,其核心设计理念是"微批次"(Micro-batch),即将实时数据流切分为小批次进行处理,兼顾了流处理的实时性和批处理的高吞吐。
DStream(离散化流)#
时间轴: t0 t1 t2 t3 t4
| | | | |
数据流: =====> =====> =====> =====> =====>
| | | | |
切分: [批次0] [批次1] [批次2] [批次3] [批次4]
| | | | |
v v v v v
RDD-0 RDD-1 RDD-2 RDD-3 RDD-4DStream(Discretized Stream)是Spark Streaming的核心抽象,它将连续的数据流按照时间间隔(batch interval)切分为一系列连续的RDD:
- DStream = RDD的时间序列:每个时间间隔内的数据形成一个RDD
- 批次间隔(Batch Interval):由
spark.streaming.batchDuration控制,通常为500ms到数秒 - 依赖关系:DStream之间的转换操作会生成RDD之间的依赖关系
窗口计算#
窗口计算是流处理中的核心能力,它允许在滑动的时间窗口内进行聚合计算:
时间轴: t0 t1 t2 t3 t4 t5 t6 t7
数据: [A] [B] [C] [D] [E] [F] [G] [H]
窗口长度=3个批次, 滑动步长=2个批次:
窗口1: [A] [B] [C] -> 计算结果1
窗口2: [C] [D] [E] -> 计算结果2
窗口3: [E] [F] [G] -> 计算结果3| 参数 | 含义 | 示例 |
|---|---|---|
| 窗口长度(Window Duration) | 窗口覆盖的时间范围 | 30秒 |
| 滑动步长(Slide Duration) | 窗口每次移动的时间间隔 | 10秒 |
| 批次间隔(Batch Interval) | 数据切分的最小单位 | 5秒 |
注意:窗口长度和滑动步长都必须是批次间隔的整数倍。
常用窗口操作:
| 操作 | 说明 | 使用场景 |
|---|---|---|
window() | 返回窗口内的原始DStream | 自定义窗口聚合 |
reduceByWindow() | 窗口内聚合 | 窗口内求和/计数 |
reduceByKeyAndWindow() | 按Key分组后窗口聚合 | 按车型统计窗口内销量 |
countByWindow() | 窗口内计数 | 窗口内消息总数 |
countByValueAndWindow() | 按值分组后窗口计数 | 窗口内各类型消息数量 |
背压机制(Backpressure)#
当数据生产速度超过消费速度时,系统需要背压机制来防止数据积压导致崩溃:
生产者 (Kafka) 消费者 (Spark Streaming)
| |
| 数据速率: 10000条/秒 | 处理速率: 5000条/秒
| ---------> | <---------
| |
| 背压控制器 |
| +----------------------+ |
| | PID控制器算法: | |
| | 1. 采集处理速率 | |
| | 2. 计算速率差值 | |
| | 3. 动态调整消费速率 | |
| +----------------------+ |
| |
| 调整后: 5000条/秒 | 处理速率: 5000条/秒
| ---------> | <---------背压机制的关键配置:
| 配置项 | 默认值 | 说明 |
|---|---|---|
spark.streaming.backpressure.enabled | false | 是否启用背压 |
spark.streaming.backpressure.initialRate | 无 | 初始最大接收速率 |
spark.streaming.receiver.maxRate | 无 | 每个接收器的最大接收速率(启用背压时的上限) |
spark.streaming.kafka.maxRatePerPartition | 无 | 每个Kafka分区的最大读取速率 |
最佳实践:在生产环境中,建议始终启用背压机制(
spark.streaming.backpressure.enabled=true),并设置合理的maxRatePerPartition作为速率上限,防止系统在突发流量时过载。
知识点4:汽车数据分析指标体系#
构建完善的指标体系是数据分析的基础。汽车行业的数据分析指标体系通常从四个维度展开:
销量分析指标#
| 指标类别 | 具体指标 | 计算方式 | 分析维度 |
|---|---|---|---|
| 销量规模 | 总销量 | SUM(销售数量) | 品牌/车型/地区/时间 |
| 销量规模 | 同比增长率 | (本期-同期)/同期 x 100% | 品牌/车型/地区 |
| 销量规模 | 环比增长率 | (本期-上期)/上期 x 100% | 品牌/车型/地区 |
| 销量结构 | 品牌市场份额 | 品牌销量/总销量 x 100% | 品牌/地区/时间 |
| 销量结构 | 车型占比 | 车型销量/品牌销量 x 100% | 品牌/车型 |
| 销量趋势 | 移动平均 | 近N期销量均值 | 品牌/车型/时间 |
价格分析指标#
| 指标类别 | 具体指标 | 计算方式 | 分析维度 |
|---|---|---|---|
| 价格水平 | 平均成交价 | AVG(成交价格) | 品牌/车型/地区 |
| 价格水平 | 中位数价格 | MEDIAN(成交价格) | 品牌/车型/地区 |
| 价格波动 | 价格离散度 | STDDEV(成交价格) | 品牌/车型 |
| 价格波动 | 优惠幅度 | (指导价-成交价)/指导价 x 100% | 品牌/车型/地区 |
| 价格弹性 | 价格弹性系数 | 销量变化率/价格变化率 | 品牌/车型 |
用户画像指标#
| 指标类别 | 具体指标 | 数据来源 | 分析维度 |
|---|---|---|---|
| 人口统计 | 年龄分布 | 客户信息 | 品牌/车型 |
| 人口统计 | 性别比例 | 客户信息 | 品牌/车型 |
| 经济能力 | 收入水平 | 客户信息 | 品牌/车型 |
| 消费偏好 | 支付方式偏好 | 交易记录 | 品牌/地区 |
| 消费偏好 | 配置偏好 | 交易记录 | 车型 |
| 行为特征 | 购车周期 | 交易记录 | 品牌/客户类型 |
竞争力分析指标#
| 指标类别 | 具体指标 | 计算方式 | 分析维度 |
|---|---|---|---|
| 市场地位 | 市场占有率 | 品牌销量/总销量 x 100% | 品牌/地区/时间 |
| 市场地位 | 市场排名 | RANK(品牌销量) | 品牌/地区 |
| 产品竞争力 | 性价比指数 | 配置得分/成交价格 | 品牌/车型 |
| 产品竞争力 | 口碑指数 | 评分加权平均 | 品牌/车型 |
| 渠道能力 | 经销商覆盖率 | 有销量的经销商数/总经销商数 | 品牌/地区 |
系统架构#
技术栈版本#
| 组件 | 版本 | 用途 |
|---|---|---|
| JDK | 11 | 运行环境 |
| Scala | 2.13.8 | 开发语言 |
| Spark | 3.5.8 | 数据处理引擎 |
| Kafka | 3.7.0 | 消息队列 |
| HBase | 2.5.9 | 时序数据存储 |
| Hadoop | 3.3.6 | 分布式存储 |
| Hive | 3.1.3 | 数据仓库 |
| MySQL | 8.0 | 业务数据存储 |
| Typesafe Config | 1.4.3 | 配置管理 |
系统架构图#
+------------------+ +------------------+ +------------------+ +------------------+
| 数据采集层 | --> | 数据存储层 | --> | 数据处理层 | --> | 可视化层 |
+------------------+ +------------------+ +------------------+ +------------------+
| - Kafka 3.7.0 | | - HDFS (Hadoop) | | - Spark Core | | - Grafana |
| - Flume | | - HBase 2.5.9 | | - Spark SQL | | - Tableau |
| - OBD数据网关 | | - MySQL 8.0 | | - Spark Streaming| | - 自研Dashboard |
+------------------+ +------------------+ | - Hive 3.1.3 | +------------------+
+------------------+
|
+------------------+
| 基础设施层 |
+------------------+
| - ConfigManager |
| - DataAccessLayer|
| - 降级策略 |
+------------------+项目工程配置#
pom.xml#
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.sparklearning</groupId>
<artifactId>auto-bigdata-analysis</artifactId>
<version>1.0.0</version>
<packaging>jar</packaging>
<name>Auto BigData Analysis System</name>
<description>汽车大数据分析系统 - 基于Spark的企业级车联网数据处理平台</description>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<java.version>11</java.version>
<scala.version>2.13.8</scala.version>
<scala.binary.version>2.13</scala.binary.version>
<spark.version>3.5.8</spark.version>
<kafka.version>3.7.0</kafka.version>
<hbase.version>2.5.9</hbase.version>
<hadoop.version>3.3.6</hadoop.version>
<hive.version>3.1.3</hive.version>
<mysql.connector.version>8.0.33</mysql.connector.version>
<typesafe.config.version>1.4.3</typesafe.config.version>
<slf4j.version>1.7.36</slf4j.version>
<log4j.version>2.17.2</log4j.version>
</properties>
<dependencies>
<!-- Scala 标准库 -->
<dependency>
<groupId>org.scala-lang</groupId>
<artifactId>scala-library</artifactId>
<version>${scala.version}</version>
</dependency>
<!-- Spark Core -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
<scope>provided</scope>
</dependency>
<!-- Spark SQL -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
<scope>provided</scope>
</dependency>
<!-- Spark Streaming -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
<scope>provided</scope>
</dependency>
<!-- Spark Streaming Kafka Connector -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming-kafka-0-10_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
</dependency>
<!-- Spark SQL Kafka Connector -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql-kafka-0-10_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
</dependency>
<!-- Kafka Client -->
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>${kafka.version}</version>
</dependency>
<!-- HBase Client -->
<dependency>
<groupId>org.apache.hbase</groupId>
<artifactId>hbase-client</artifactId>
<version>${hbase.version}</version>
</dependency>
<!-- HBase MapReduce (用于Spark与HBase集成) -->
<dependency>
<groupId>org.apache.hbase</groupId>
<artifactId>hbase-mapreduce</artifactId>
<version>${hbase.version}</version>
</dependency>
<!-- Hadoop Common -->
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-common</artifactId>
<version>${hadoop.version}</version>
<scope>provided</scope>
</dependency>
<!-- Hadoop HDFS -->
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-hdfs</artifactId>
<version>${hadoop.version}</version>
<scope>provided</scope>
</dependency>
<!-- Hive JDBC -->
<dependency>
<groupId>org.apache.hive</groupId>
<artifactId>hive-jdbc</artifactId>
<version>${hive.version}</version>
<scope>provided</scope>
</dependency>
<!-- MySQL Connector -->
<dependency>
<groupId>com.mysql</groupId>
<artifactId>mysql-connector-j</artifactId>
<version>${mysql.connector.version}</version>
</dependency>
<!-- Typesafe Config -->
<dependency>
<groupId>com.typesafe</groupId>
<artifactId>config</artifactId>
<version>${typesafe.config.version}</version>
</dependency>
<!-- SLF4J API -->
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
<version>${slf4j.version}</version>
</dependency>
<!-- Log4j2 -->
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-core</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
<dependency>
<groupId>org.apache.logging.log4j</groupId>
<artifactId>log4j-slf4j-impl</artifactId>
<version>${log4j.version}</version>
<scope>runtime</scope>
</dependency>
</dependencies>
<build>
<plugins>
<!-- Scala Maven Plugin -->
<plugin>
<groupId>net.alchim31.maven</groupId>
<artifactId>scala-maven-plugin</artifactId>
<version>4.8.1</version>
<executions>
<execution>
<id>scala-compile-first</id>
<phase>process-resources</phase>
<goals>
<goal>add-source</goal>
<goal>compile</goal>
</goals>
</execution>
<execution>
<id>scala-test-compile</id>
<phase>process-test-resources</phase>
<goals>
<goal>testCompile</goal>
</goals>
</execution>
</executions>
<configuration>
<scalaVersion>${scala.version}</scalaVersion>
<args>
<arg>-target:jvm-11</arg>
<arg>-deprecation</arg>
<arg>-feature</arg>
<arg>-unchecked</arg>
</args>
<jvmArgs>
<jvmArg>-Xms128m</jvmArg>
<jvmArg>-Xmx1024m</jvmArg>
</jvmArgs>
</configuration>
</plugin>
<!-- Maven Compiler Plugin -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.11.0</version>
<configuration>
<source>${java.version}</source>
<target>${java.version}</target>
<encoding>UTF-8</encoding>
</configuration>
</plugin>
<!-- Maven Shade Plugin (构建Fat Jar) -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.5.1</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<filters>
<filter>
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
<transformers>
<transformer implementation="org.apache.maven.plugins.shade.resource.ServicesResourceTransformer"/>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>application.conf 配置文件#
将所有配置外部化到application.conf中,避免在代码中硬编码任何连接信息或密码。该文件放置在src/main/resources/目录下。
# ========================================
# 汽车大数据分析系统 - 应用配置文件
# ========================================
# 注意:生产环境中密码应通过环境变量或密钥管理系统注入
# 此处仅展示配置结构,实际部署时请替换为真实值
# ----------------------------------------
# MySQL 数据库配置
# ----------------------------------------
mysql {
url = "jdbc:mysql://localhost:3306/auto_database?useSSL=true&serverTimezone=Asia/Shanghai&characterEncoding=utf8mb4"
url = ${?MYSQL_URL}
user = "auto_admin"
user = ${?MYSQL_USER}
password = "请通过环境变量MYSQL_PASSWORD设置"
password = ${?MYSQL_PASSWORD}
driver = "com.mysql.cj.jdbc.Driver"
connection-pool {
initial-size = 5
max-active = 20
max-idle = 10
min-idle = 5
max-wait-ms = 10000
}
}
# ----------------------------------------
# Kafka 配置
# ----------------------------------------
kafka {
bootstrap-servers = "localhost:9092"
bootstrap-servers = ${?KAFKA_BOOTSTRAP_SERVERS}
# 消费者配置
consumer {
group-id = "auto-bigdata-consumer-group"
auto-offset-reset = "latest"
enable-auto-commit = false
max-poll-records = 500
session-timeout-ms = 30000
heartbeat-interval-ms = 10000
}
# 生产者配置
producer {
acks = "all"
retries = 3
batch-size = 16384
linger-ms = 5
buffer-memory = 33554432
}
# 主题配置
topics {
sales-data = "sales-data"
sensor-data = "sensor-data"
gps-data = "gps-data"
obd-data = "obd-data"
alert-data = "alert-data"
}
}
# ----------------------------------------
# HBase 配置
# ----------------------------------------
hbase {
zookeeper-quorum = "localhost"
zookeeper-quorum = ${?HBASE_ZOOKEEPER_QUORUM}
zookeeper-client-port = "2181"
zookeeper-client-port = ${?HBASE_ZOOKEEPER_CLIENT_PORT}
rpc-timeout = 60000
client-retries-number = 3
client-operation-timeout = 1200000
# 表命名空间
tables {
sensor-data = "sensor_data"
gps-track = "gps_track"
obd-diagnosis = "obd_diagnosis"
driving-behavior = "driving_behavior"
}
# 列族配置
column-families {
default-cf = "cf"
detail-cf = "detail"
}
}
# ----------------------------------------
# Spark 配置
# ----------------------------------------
spark {
app-name = "AutoBigDataAnalysis"
master = "local[*]"
master = ${?SPARK_MASTER}
# Spark SQL配置
sql {
shuffle-partitions = 200
adaptive-execution-enabled = true
warehouse-dir = "/user/hive/warehouse"
}
# Spark Streaming配置
streaming {
batch-duration-seconds = 5
backpressure-enabled = true
backpressure-initial-rate = 1000
kafka-max-rate-per-partition = 100
checkpoint-dir = "/tmp/spark-checkpoint/auto-bigdata"
}
# 序列化配置
serializer = "org.apache.spark.serializer.KryoSerializer"
kryo-classes = [
"org.apache.hadoop.hbase.client.Result",
"org.apache.hadoop.hbase.io.ImmutableBytesWritable"
]
}
# ----------------------------------------
# 数据验证配置
# ----------------------------------------
validation {
# 传感器值范围
sensor-ranges {
"发动机温度" { min = -40.0, max = 150.0 }
"机油压力" { min = 0.0, max = 10.0 }
"冷却液温度" { min = -40.0, max = 130.0 }
"轮胎压力" { min = 0.0, max = 5.0 }
"燃油液位" { min = 0.0, max = 100.0 }
"电池电压" { min = 0.0, max = 16.0 }
}
# 销售价格范围
sale-price {
min = 10000.0
max = 10000000.0
}
# 客户年龄范围
customer-age {
min = 18
max = 100
}
}
# ----------------------------------------
# 降级策略配置
# ----------------------------------------
degradation {
# 实时处理失败时的回退方案
enabled = true
# 最大重试次数
max-retries = 3
# 重试间隔(毫秒)
retry-interval-ms = 5000
# 降级到批处理的阈值: 连续失败次数
batch-fallback-threshold = 5
# 降级时数据暂存路径
fallback-data-path = "/tmp/auto-bigdata/fallback"
# 告警通知配置
alert {
enabled = true
email = "admin@example.com"
email = ${?ALERT_EMAIL}
}
}
# ----------------------------------------
# 数据库表名配置
# ----------------------------------------
database {
tables {
sales = "sales"
vehicles = "vehicles"
customers = "customers"
region-sales-analysis = "region_sales_analysis"
payment-sales-analysis = "payment_sales_analysis"
sales-trend-analysis = "sales_trend_analysis"
sensor-abnormal-analysis = "sensor_abnormal_analysis"
sensor-avg-analysis = "sensor_avg_analysis"
vehicle-competitiveness = "vehicle_competitiveness"
customer-portrait = "customer_portrait"
}
}核心基础设施代码#
ConfigManager.scala - 配置管理器#
使用Typesafe Config库实现统一的配置管理,所有配置项从application.conf读取,敏感信息通过环境变量覆盖。
package com.sparklearning.auto.config
import com.typesafe.config.{Config, ConfigFactory}
import org.slf4j.{Logger, LoggerFactory}
import scala.jdk.CollectionConverters._
import scala.util.{Try, Success, Failure}
/**
* 配置管理器 - 统一管理应用配置
*
* 设计原则:
* 1. 所有配置从application.conf读取,禁止硬编码
* 2. 敏感信息(密码等)通过环境变量覆盖
* 3. 提供类型安全的配置访问方法
* 4. 配置加载失败时提供有意义的错误信息
*
* 使用方式:
* val config = ConfigManager.getConfig
* val mysqlUrl = ConfigManager.getMySQLUrl
*/
object ConfigManager {
private val logger: Logger = LoggerFactory.getLogger(ConfigManager.getClass)
// 加载配置,优先加载环境变量覆盖
private lazy val config: Config = {
val baseConfig = ConfigFactory.load()
val envConfig = ConfigFactory.systemEnvironment()
// 环境变量优先级最高,然后是系统属性,最后是配置文件
ConfigFactory.load(
ConfigFactory.systemProperties()
.withFallback(envConfig)
.withFallback(baseConfig)
).resolve()
}
/**
* 获取完整配置对象
*/
def getConfig: Config = config
// ========================================
// MySQL 配置
// ========================================
def getMySQLUrl: String = config.getString("mysql.url")
def getMySQLUser: String = config.getString("mysql.user")
def getMySQLPassword: String = {
val password = config.getString("mysql.password")
if (password.contains("请通过环境变量") || password.isEmpty) {
logger.warn("MySQL密码未通过环境变量配置,请在MYSQL_PASSWORD环境变量中设置密码")
}
password
}
def getMySQLDriver: String = config.getString("mysql.driver")
/**
* 获取MySQL连接选项(用于Spark JDBC)
*/
def getMySQLOptions(table: String): Map[String, String] = {
require(table.nonEmpty, "表名不能为空")
Map(
"url" -> getMySQLUrl,
"dbtable" -> table,
"user" -> getMySQLUser,
"password" -> getMySQLPassword,
"driver" -> getMySQLDriver
)
}
// ========================================
// Kafka 配置
// ========================================
def getKafkaBootstrapServers: String = config.getString("kafka.bootstrap-servers")
def getKafkaConsumerGroupId: String = config.getString("kafka.consumer.group-id")
def getKafkaConsumerConfig: Map[String, String] = Map(
"bootstrap.servers" -> getKafkaBootstrapServers,
"group.id" -> getKafkaConsumerGroupId,
"auto.offset.reset" -> config.getString("kafka.consumer.auto-offset-reset"),
"enable.auto.commit" -> config.getString("kafka.consumer.enable-auto-commit"),
"max.poll.records" -> config.getString("kafka.consumer.max-poll-records"),
"session.timeout.ms" -> config.getString("kafka.consumer.session-timeout-ms"),
"heartbeat.interval.ms" -> config.getString("kafka.consumer.heartbeat-interval-ms")
)
def getKafkaProducerConfig: java.util.Properties = {
val props = new java.util.Properties()
props.put("bootstrap.servers", getKafkaBootstrapServers)
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
props.put("acks", config.getString("kafka.producer.acks"))
props.put("retries", config.getString("kafka.producer.retries"))
props.put("batch.size", config.getString("kafka.producer.batch-size"))
props.put("linger.ms", config.getString("kafka.producer.linger-ms"))
props.put("buffer.memory", config.getString("kafka.producer.buffer-memory"))
props
}
def getKafkaTopic(topicKey: String): String = {
Try(config.getString(s"kafka.topics.$topicKey")) match {
case Success(topic) => topic
case Failure(_) =>
logger.warn(s"未找到Kafka主题配置: kafka.topics.$topicKey,使用默认值: $topicKey")
topicKey
}
}
// ========================================
// HBase 配置
// ========================================
def getHBaseZookeeperQuorum: String = config.getString("hbase.zookeeper-quorum")
def getHBaseZookeeperClientPort: String = config.getString("hbase.zookeeper-client-port")
def getHBaseTableName(tableKey: String): String = {
Try(config.getString(s"hbase.tables.$tableKey")) match {
case Success(table) => table
case Failure(_) =>
logger.warn(s"未找到HBase表配置: hbase.tables.$tableKey,使用默认值: $tableKey")
tableKey
}
}
/**
* 获取HBase Configuration对象
*/
def getHBaseConfig: org.apache.hadoop.conf.Configuration = {
val conf = org.apache.hadoop.hbase.HBaseConfiguration.create()
conf.set("hbase.zookeeper.quorum", getHBaseZookeeperQuorum)
conf.set("hbase.zookeeper.property.clientPort", getHBaseZookeeperClientPort)
conf.set("hbase.rpc.timeout", config.getString("hbase.rpc-timeout").toString)
conf.set("hbase.client.retries.number", config.getString("hbase.client-retries-number").toString)
conf.set("hbase.client.operation.timeout", config.getString("hbase.client-operation-timeout").toString)
conf
}
// ========================================
// Spark 配置
// ========================================
def getSparkAppName: String = config.getString("spark.app-name")
def getSparkMaster: String = config.getString("spark.master")
def getSparkBatchDurationSeconds: Int = config.getInt("spark.streaming.batch-duration-seconds")
def getSparkCheckpointDir: String = config.getString("spark.streaming.checkpoint-dir")
def getSparkShufflePartitions: Int = config.getInt("spark.sql.shuffle-partitions")
// ========================================
// 数据验证配置
// ========================================
/**
* 获取传感器值范围
* @param sensorType 传感器类型
* @return (最小值, 最大值)
*/
def getSensorRange(sensorType: String): (Double, Double) = {
val path = s"validation.sensor-ranges.$sensorType"
Try {
val min = config.getDouble(s"$path.min")
val max = config.getDouble(s"$path.max")
(min, max)
} match {
case Success(range) => range
case Failure(_) =>
logger.warn(s"未找到传感器范围配置: $path,使用默认范围(0.0, 100.0)")
(0.0, 100.0)
}
}
def getSalePriceRange: (Double, Double) = {
(config.getDouble("validation.sale-price.min"), config.getDouble("validation.sale-price.max"))
}
def getCustomerAgeRange: (Int, Int) = {
(config.getInt("validation.customer-age.min"), config.getInt("validation.customer-age.max"))
}
// ========================================
// 降级策略配置
// ========================================
def isDegradationEnabled: Boolean = config.getBoolean("degradation.enabled")
def getMaxRetries: Int = config.getInt("degradation.max-retries")
def getRetryIntervalMs: Long = config.getLong("degradation.retry-interval-ms")
def getBatchFallbackThreshold: Int = config.getInt("degradation.batch-fallback-threshold")
def getFallbackDataPath: String = config.getString("degradation.fallback-data-path")
// ========================================
// 数据库表名配置
// ========================================
def getDatabaseTable(tableKey: String): String = {
Try(config.getString(s"database.tables.$tableKey")) match {
case Success(table) => table
case Failure(_) =>
logger.warn(s"未找到数据库表配置: database.tables.$tableKey,使用默认值: $tableKey")
tableKey
}
}
}DataAccessLayer.scala - 数据访问层#
统一封装MySQL和HBase的数据访问操作,提供连接管理、资源释放和异常处理。
package com.sparklearning.auto.dataaccess
import com.sparklearning.auto.config.ConfigManager
import org.apache.hadoop.hbase.{HBaseConfiguration, TableName}
import org.apache.hadoop.hbase.client.{Connection, ConnectionFactory, Put, Scan, Table}
import org.apache.hadoop.hbase.util.Bytes
import org.slf4j.{Logger, LoggerFactory}
import java.sql.{Connection => SQLConnection, DriverManager, PreparedStatement, ResultSet}
import scala.jdk.CollectionConverters._
import scala.util.{Try, Success, Failure, Using}
/**
* 数据访问层 - 统一封装数据库操作
*
* 设计原则:
* 1. 所有数据库连接通过ConfigManager获取配置
* 2. 使用loan pattern确保资源释放
* 3. 提供统一的异常处理和日志记录
* 4. 连接池管理(简化实现,生产环境建议使用HikariCP)
*/
object DataAccessLayer {
private val logger: Logger = LoggerFactory.getLogger(DataAccessLayer.getClass)
// ========================================
// MySQL 数据访问
// ========================================
/**
* 获取MySQL连接
* 使用loan pattern,确保连接在使用后自动关闭
*
* @param operation 数据库操作函数
* @tparam T 返回值类型
* @return 操作结果
*/
def withMySQLConnection[T](operation: SQLConnection => T): Try[T] = {
var connection: SQLConnection = null
try {
Class.forName(ConfigManager.getMySQLDriver)
connection = DriverManager.getConnection(
ConfigManager.getMySQLUrl,
ConfigManager.getMySQLUser,
ConfigManager.getMySQLPassword
)
connection.setAutoCommit(false)
val result = operation(connection)
connection.commit()
Success(result)
} catch {
case e: Exception =>
if (connection != null) {
try {
connection.rollback()
} catch {
case rollbackEx: Exception =>
logger.error("MySQL回滚失败", rollbackEx)
}
}
logger.error("MySQL操作失败", e)
Failure(e)
} finally {
if (connection != null) {
try {
connection.close()
} catch {
case e: Exception =>
logger.error("MySQL连接关闭失败", e)
}
}
}
}
/**
* 批量写入MySQL
*
* @param tableKey 表配置键名
* @param columns 列名列表
* @param records 数据记录列表
* @return 成功写入的记录数
*/
def batchInsertToMySQL(tableKey: String, columns: Seq[String],
records: Seq[Seq[Any]]): Try[Int] = {
withMySQLConnection { conn =>
val tableName = ConfigManager.getDatabaseTable(tableKey)
val placeholders = columns.map(_ => "?").mkString(",")
val sql = s"INSERT INTO $tableName (${columns.mkString(",")}) VALUES ($placeholders)"
var successCount = 0
val stmt = conn.prepareStatement(sql)
try {
records.foreach { record =>
record.zipWithIndex.foreach { case (value, idx) =>
stmt.setObject(idx + 1, value)
}
stmt.addBatch()
successCount += 1
}
stmt.executeBatch()
successCount
} finally {
stmt.close()
}
}
}
/**
* 执行查询并处理结果
*
* @param sql SQL查询语句
* @param processor 结果处理函数
* @tparam T 返回值类型
* @return 查询结果
*/
def queryMySQL[T](sql: String)(processor: ResultSet => T): Try[T] = {
withMySQLConnection { conn =>
val stmt = conn.createStatement()
var rs: ResultSet = null
try {
rs = stmt.executeQuery(sql)
processor(rs)
} finally {
if (rs != null) rs.close()
stmt.close()
}
}
}
// ========================================
// HBase 数据访问
// ========================================
/**
* 获取HBase连接
* 使用loan pattern,确保连接在使用后自动关闭
*
* @param operation HBase操作函数
* @tparam T 返回值类型
* @return 操作结果
*/
def withHBaseConnection[T](operation: Connection => T): Try[T] = {
var connection: Connection = null
try {
connection = ConnectionFactory.createConnection(ConfigManager.getHBaseConfig)
Success(operation(connection))
} catch {
case e: Exception =>
logger.error("HBase操作失败", e)
Failure(e)
} finally {
if (connection != null) {
try {
connection.close()
} catch {
case e: Exception =>
logger.error("HBase连接关闭失败", e)
}
}
}
}
/**
* 批量写入HBase
*
* @param tableKey 表配置键名
* @param columnFamily 列族名
* @param records 数据记录列表: (rowKey, Seq[(columnQualifier, value)])
* @return 成功写入的记录数
*/
def batchInsertToHBase(tableKey: String, columnFamily: String,
records: Seq[(String, Seq[(String, Any)])]): Try[Int] = {
withHBaseConnection { conn =>
val tableName = ConfigManager.getHBaseTableName(tableKey)
val table = conn.getTable(TableName.valueOf(tableName))
try {
val cfBytes = Bytes.toBytes(columnFamily)
val puts = records.map { case (rowKey, columns) =>
val put = new Put(Bytes.toBytes(rowKey))
columns.foreach { case (qualifier, value) =>
value match {
case v: Int => put.addColumn(cfBytes, Bytes.toBytes(qualifier), Bytes.toBytes(v))
case v: Long => put.addColumn(cfBytes, Bytes.toBytes(qualifier), Bytes.toBytes(v))
case v: Double => put.addColumn(cfBytes, Bytes.toBytes(qualifier), Bytes.toBytes(v))
case v: String => put.addColumn(cfBytes, Bytes.toBytes(qualifier), Bytes.toBytes(v))
case v: Float => put.addColumn(cfBytes, Bytes.toBytes(qualifier), Bytes.toBytes(v))
case v: Boolean => put.addColumn(cfBytes, Bytes.toBytes(qualifier), Bytes.toBytes(v))
case v: Array[Byte] => put.addColumn(cfBytes, Bytes.toBytes(qualifier), v)
case other =>
put.addColumn(cfBytes, Bytes.toBytes(qualifier), Bytes.toBytes(other.toString))
}
}
put
}
table.put(puts.asJava)
records.size
} finally {
table.close()
}
}
}
// ========================================
// 数据验证工具
// ========================================
/**
* 验证传感器值是否在合理范围内
*/
def validateSensorValue(sensorType: String, value: Double): Boolean = {
val (min, max) = ConfigManager.getSensorRange(sensorType)
if (value < min || value > max) {
logger.warn(s"传感器值超出合理范围: type=$sensorType, value=$value, range=[$min, $max]")
false
} else {
true
}
}
/**
* 验证销售价格是否在合理范围内
*/
def validateSalePrice(price: Double): Boolean = {
val (min, max) = ConfigManager.getSalePriceRange
if (price < min || price > max) {
logger.warn(s"销售价格超出合理范围: price=$price, range=[$min, $max]")
false
} else {
true
}
}
/**
* 验证客户年龄是否在合理范围内
*/
def validateCustomerAge(age: Int): Boolean = {
val (min, max) = ConfigManager.getCustomerAgeRange
if (age < min || age > max) {
logger.warn(s"客户年龄超出合理范围: age=$age, range=[$min, $max]")
false
} else {
true
}
}
/**
* 通用空值检查
*/
def requireNonEmpty(value: String, fieldName: String): Unit = {
require(value != null && value.trim.nonEmpty, s"$fieldName 不能为空")
}
}数据模型设计#
1. 销售数据模型#
| 字段名 | 数据类型 | 描述 | 约束 |
|---|---|---|---|
| sale_id | BIGINT | 销售记录ID | PRIMARY KEY |
| vehicle_id | INT | 车辆ID | NOT NULL |
| customer_id | INT | 客户ID | NOT NULL |
| dealer_id | INT | 经销商ID | NOT NULL |
| sale_date | DATE | 销售日期 | NOT NULL |
| sale_price | DECIMAL(12,2) | 销售价格 | CHECK(sale_price > 0) |
| payment_method | VARCHAR(50) | 支付方式 | NOT NULL |
| region | VARCHAR(100) | 销售地区 | NOT NULL |
2. 车辆数据模型#
| 字段名 | 数据类型 | 描述 | 约束 |
|---|---|---|---|
| vehicle_id | INT | 车辆ID | PRIMARY KEY |
| brand | VARCHAR(50) | 品牌 | NOT NULL |
| model | VARCHAR(100) | 型号 | NOT NULL |
| year | INT | 年份 | CHECK(year >= 2000 AND year <= 2030) |
| fuel_type | VARCHAR(20) | 燃油类型 | NOT NULL |
| engine_size | DECIMAL(3,1) | 发动机排量 | CHECK(engine_size > 0) |
| horsepower | INT | 马力 | CHECK(horsepower > 0) |
| price | DECIMAL(12,2) | 指导价格 | CHECK(price > 0) |
| features | JSON | 配置特征 |
3. 客户数据模型#
| 字段名 | 数据类型 | 描述 | 约束 |
|---|---|---|---|
| customer_id | INT | 客户ID | PRIMARY KEY |
| name | VARCHAR(100) | 姓名 | NOT NULL |
| age | INT | 年龄 | CHECK(age >= 18 AND age <= 100) |
| gender | VARCHAR(10) | 性别 | NOT NULL |
| occupation | VARCHAR(100) | 职业 | |
| income | DECIMAL(12,2) | 收入 | CHECK(income >= 0) |
| region | VARCHAR(100) | 地区 | NOT NULL |
| purchase_history | JSON | 购买历史 |
4. 传感器数据模型#
| 字段名 | 数据类型 | 描述 | 约束 |
|---|---|---|---|
| sensor_id | INT | 传感器ID | NOT NULL |
| vehicle_id | INT | 车辆ID | NOT NULL |
| timestamp | BIGINT | 采集时间(毫秒时间戳) | NOT NULL |
| sensor_type | VARCHAR(50) | 传感器类型 | NOT NULL |
| value | DOUBLE | 传感器值 | NOT NULL |
| status | VARCHAR(20) | 状态 | NOT NULL |
5. GPS轨迹数据模型#
| 字段名 | 数据类型 | 描述 | 约束 |
|---|---|---|---|
| vehicle_id | INT | 车辆ID | NOT NULL |
| timestamp | BIGINT | 采集时间(毫秒时间戳) | NOT NULL |
| latitude | DOUBLE | 纬度 | CHECK(latitude BETWEEN -90 AND 90) |
| longitude | DOUBLE | 经度 | CHECK(longitude BETWEEN -180 AND 180) |
| speed | DOUBLE | 速度(km/h) | CHECK(speed >= 0) |
| heading | DOUBLE | 航向角(0-360) |
系统模块设计#
1. 数据采集模块#
1.1 销售数据采集#
使用Kafka采集销售系统的实时数据,所有配置通过ConfigManager获取:
package com.sparklearning.auto.producer
import com.sparklearning.auto.config.ConfigManager
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord}
import org.slf4j.{Logger, LoggerFactory}
import scala.util.Random
/**
* 销售数据生产者 - 模拟销售数据发送到Kafka
*
* 企业级改进:
* 1. 配置从ConfigManager获取,不硬编码
* 2. 使用try-finally确保Producer资源释放
* 3. 添加数据验证逻辑
* 4. 添加日志记录
*/
object SalesDataProducer {
private val logger: Logger = LoggerFactory.getLogger(SalesDataProducer.getClass)
def main(args: Array[String]): Unit = {
val producerProps = ConfigManager.getKafkaProducerConfig
val topic = ConfigManager.getKafkaTopic("sales-data")
val producer = new KafkaProducer[String, String](producerProps)
val random = new Random()
val brands = Array("比亚迪", "吉利", "长安", "奇瑞", "哈弗", "红旗",
"蔚来", "理想", "小鹏", "问界", "极氪", "领克")
val models = Array("秦PLUS", "星瑞", "UNI-V", "瑞虎8", "H6", "H5",
"ES6", "L7", "P7", "M5", "001", "08")
val regions = Array("北京", "上海", "广州", "深圳", "杭州", "成都", "武汉", "西安",
"南京", "重庆", "苏州", "天津")
val paymentMethods = Array("现金", "贷款", "分期付款", "融资租赁")
try {
var sentCount = 0L
while (true) {
val saleId = System.currentTimeMillis()
val vehicleId = 1000 + random.nextInt(9000)
val customerId = 100 + random.nextInt(900)
val dealerId = 10 + random.nextInt(90)
val year = 2023 + random.nextInt(4)
val month = 1 + random.nextInt(12)
val day = 1 + random.nextInt(28)
val saleDate = f"$year-$month%02d-$day%02d"
val salePrice = 100000.0 + random.nextInt(900000)
val paymentMethod = paymentMethods(random.nextInt(paymentMethods.length))
val region = regions(random.nextInt(regions.length))
// 数据验证
if (!DataValidator.validateSalePrice(salePrice)) {
logger.warn(s"跳过无效销售价格: $salePrice")
Thread.sleep(2000)
next
}
val salesData = s"""{"sale_id":$saleId,"vehicle_id":$vehicleId,"customer_id":$customerId,"dealer_id":$dealerId,"sale_date":"$saleDate","sale_price":$salePrice,"payment_method":"$paymentMethod","region":"$region"}"""
val record = new ProducerRecord[String, String](topic, null, salesData)
producer.send(record)
sentCount += 1
if (sentCount % 100 == 0) {
logger.info(s"已发送销售数据: $sentCount 条")
}
Thread.sleep(2000)
}
} catch {
case e: InterruptedException =>
logger.info("销售数据生产者被中断,正在关闭...")
case e: Exception =>
logger.error("销售数据生产者异常", e)
} finally {
producer.close()
logger.info("Kafka Producer已关闭")
}
}
}
/**
* 数据验证工具 - 集中管理数据验证逻辑
*/
object DataValidator {
private val logger: Logger = LoggerFactory.getLogger(DataValidator.getClass)
/**
* 验证销售价格
*/
def validateSalePrice(price: Double): Boolean = {
val (min, max) = ConfigManager.getSalePriceRange
val valid = price >= min && price <= max
if (!valid) {
logger.warn(s"销售价格超出合理范围: price=$price, range=[$min, $max]")
}
valid
}
/**
* 验证传感器值
*/
def validateSensorValue(sensorType: String, value: Double): Boolean = {
val (min, max) = ConfigManager.getSensorRange(sensorType)
val valid = value >= min && value <= max
if (!valid) {
logger.warn(s"传感器值超出合理范围: type=$sensorType, value=$value, range=[$min, $max]")
}
valid
}
/**
* 验证非空字符串
*/
def validateNonEmpty(value: String, fieldName: String): Boolean = {
val valid = value != null && value.trim.nonEmpty
if (!valid) {
logger.warn(s"$fieldName 值为空")
}
valid
}
/**
* 验证数值范围
*/
def validateRange(value: Double, min: Double, max: Double, fieldName: String): Boolean = {
val valid = value >= min && value <= max
if (!valid) {
logger.warn(s"$fieldName 超出范围: value=$value, range=[$min, $max]")
}
valid
}
}1.2 传感器数据采集#
package com.sparklearning.auto.producer
import com.sparklearning.auto.config.ConfigManager
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord}
import org.slf4j.{Logger, LoggerFactory}
import scala.util.Random
/**
* 传感器数据生产者 - 模拟车辆传感器数据发送到Kafka
*
* 企业级改进:
* 1. 配置从ConfigManager获取
* 2. 传感器值范围验证
* 3. try-finally确保资源释放
* 4. 完善的异常处理
*/
object SensorDataProducer {
private val logger: Logger = LoggerFactory.getLogger(SensorDataProducer.getClass)
def main(args: Array[String]): Unit = {
val producerProps = ConfigManager.getKafkaProducerConfig
val topic = ConfigManager.getKafkaTopic("sensor-data")
val producer = new KafkaProducer[String, String](producerProps)
val random = new Random()
val sensorTypes = Array("发动机温度", "机油压力", "冷却液温度", "轮胎压力", "燃油液位", "电池电压")
// 各传感器的正常值范围
val sensorNormalRanges = Map(
"发动机温度" -> (70.0, 105.0),
"机油压力" -> (1.5, 5.0),
"冷却液温度" -> (75.0, 100.0),
"轮胎压力" -> (2.0, 2.8),
"燃油液位" -> (5.0, 95.0),
"电池电压" -> (11.5, 14.5)
)
try {
var sentCount = 0L
while (true) {
val sensorId = 1000 + random.nextInt(9000)
val vehicleId = 1000 + random.nextInt(9000)
val timestamp = System.currentTimeMillis()
val sensorType = sensorTypes(random.nextInt(sensorTypes.length))
val (normalMin, normalMax) = sensorNormalRanges(sensorType)
// 90%概率生成正常值,10%概率生成异常值
val value = if (random.nextDouble() < 0.9) {
normalMin + random.nextDouble() * (normalMax - normalMin)
} else {
// 生成异常值:超出正常范围
val (globalMin, globalMax) = ConfigManager.getSensorRange(sensorType)
if (random.nextBoolean()) {
globalMin + random.nextDouble() * (normalMin - globalMin)
} else {
normalMax + random.nextDouble() * (globalMax - normalMax)
}
}
val roundedValue = BigDecimal(value).setScale(1, BigDecimal.RoundingMode.HALF_UP).toDouble
val status = if (roundedValue < normalMin || roundedValue > normalMax) "异常" else "正常"
val sensorData = s"""{"sensor_id":$sensorId,"vehicle_id":$vehicleId,"timestamp":$timestamp,"sensor_type":"$sensorType","value":$roundedValue,"status":"$status"}"""
val record = new ProducerRecord[String, String](topic, String.valueOf(vehicleId), sensorData)
producer.send(record)
sentCount += 1
if (sentCount % 100 == 0) {
logger.info(s"已发送传感器数据: $sentCount 条")
}
Thread.sleep(1000)
}
} catch {
case e: InterruptedException =>
logger.info("传感器数据生产者被中断,正在关闭...")
case e: Exception =>
logger.error("传感器数据生产者异常", e)
} finally {
producer.close()
logger.info("Kafka Producer已关闭")
}
}
}2. 数据存储模块#
2.1 HBase存储传感器数据#
package com.sparklearning.auto.storage
import com.sparklearning.auto.config.ConfigManager
import com.sparklearning.auto.dataaccess.DataAccessLayer
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.OutputMode
import org.apache.spark.sql.types._
import org.slf4j.{Logger, LoggerFactory}
/**
* 传感器数据写入HBase
*
* 企业级改进:
* 1. SparkSession使用try-finally确保释放
* 2. HBase连接配置从ConfigManager获取
* 3. 数据验证:过滤无效传感器值
* 4. 使用foreachBatch + DataAccessLayer写入HBase
* 5. 添加降级策略:HBase不可用时暂存到HDFS
*/
object SensorDataToHBase {
private val logger: Logger = LoggerFactory.getLogger(SensorDataToHBase.getClass)
def main(args: Array[String]): Unit = {
try {
val spark = SparkSession.builder()
.appName("SensorDataToHBase")
.master(ConfigManager.getSparkMaster)
.config("spark.sql.shuffle.partitions", ConfigManager.getSparkShufflePartitions)
.getOrCreate()
import spark.implicits._
// 从Kafka读取传感器数据
val sensorStream = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", ConfigManager.getKafkaBootstrapServers)
.option("subscribe", ConfigManager.getKafkaTopic("sensor-data"))
.option("startingOffsets", "latest")
.option("maxOffsetsPerTrigger", "1000")
.load()
// 定义Schema
val schema = new StructType()
.add("sensor_id", IntegerType)
.add("vehicle_id", IntegerType)
.add("timestamp", LongType)
.add("sensor_type", StringType)
.add("value", DoubleType)
.add("status", StringType)
// 解析JSON数据 + 数据验证
val sensorData = sensorStream
.selectExpr("CAST(value AS STRING)")
.select(from_json($"value", schema).as("data"))
.select("data.*")
.filter($"data".isNotNull) // 过滤解析失败的数据
.filter($"sensor_id".isNotNull && $"vehicle_id".isNotNull && $"timestamp".isNotNull)
.filter($"sensor_type".isNotNull && $"value".isNotNull && $"status".isNotNull)
// 写入HBase(带降级策略)
val query = sensorData.writeStream
.foreachBatch { (batchDF: org.apache.spark.sql.DataFrame, batchId: Long) =>
logger.info(s"处理批次: $batchId, 记录数: ${batchDF.count()}")
try {
// 尝试写入HBase
writeBatchToHBase(batchDF, batchId)
} catch {
case e: Exception =>
logger.error(s"批次 $batchId 写入HBase失败,启动降级策略", e)
// 降级:将数据暂存到HDFS
fallbackToHDFS(batchDF, batchId)
}
}
.outputMode(OutputMode.Update())
.option("checkpointLocation", ConfigManager.getSparkCheckpointDir + "/sensor-to-hbase")
.start()
query.awaitTermination()
} catch {
case e: Exception =>
logger.error("SensorDataToHBase运行异常", e)
throw e
} finally {
val s = SparkSession.getActiveSession.orNull
if (s != null) {
s.stop()
logger.info("SparkSession已关闭")
}
}
}
/**
* 将批次数据写入HBase
*/
private def writeBatchToHBase(batchDF: org.apache.spark.sql.DataFrame, batchId: Long): Unit = {
val records = batchDF.collect().flatMap { row =>
val sensorType = row.getAs[String]("sensor_type")
val value = row.getAs[Double]("value")
// 数据验证
if (!DataAccessLayer.validateSensorValue(sensorType, value)) {
logger.warn(s"批次 $batchId: 跳过无效传感器数据, type=$sensorType, value=$value")
None
} else {
val rowKey = s"${row.getAs[Long]("timestamp")}-${row.getAs[Int]("sensor_id")}"
val columns = Seq(
("vehicle_id", row.getAs[Int]("vehicle_id").asInstanceOf[Any]),
("sensor_type", sensorType.asInstanceOf[Any]),
("value", value.asInstanceOf[Any]),
("status", row.getAs[String]("status").asInstanceOf[Any])
)
Some((rowKey, columns))
}
}
if (records.nonEmpty) {
val result = DataAccessLayer.batchInsertToHBase("sensor-data", "cf", records)
result match {
case Success(count) => logger.info(s"批次 $batchId 成功写入HBase: $count 条记录")
case Failure(e) => throw e
}
}
}
/**
* 降级策略:将数据暂存到HDFS
*/
private def fallbackToHDFS(batchDF: org.apache.spark.sql.DataFrame, batchId: Long): Unit = {
if (!ConfigManager.isDegradationEnabled) {
logger.warn("降级策略未启用,数据可能丢失")
return
}
val fallbackPath = s"${ConfigManager.getFallbackDataPath}/sensor_data/batch_$batchId"
logger.info(s"降级策略:将批次 $batchId 数据暂存到 $fallbackPath")
try {
batchDF.write
.mode("overwrite")
.json(fallbackPath)
logger.info(s"批次 $batchId 数据已暂存到HDFS: $fallbackPath")
} catch {
case e: Exception =>
logger.error(s"降级策略执行失败,批次 $batchId 数据可能丢失", e)
}
}
}2.2 MySQL存储销售数据#
package com.sparklearning.auto.storage
import com.sparklearning.auto.config.ConfigManager
import com.sparklearning.auto.dataaccess.DataAccessLayer
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.OutputMode
import org.apache.spark.sql.types._
import org.slf4j.{Logger, LoggerFactory}
/**
* 销售数据写入MySQL
*
* 企业级改进:
* 1. SparkSession使用try-finally确保释放
* 2. MySQL连接配置从ConfigManager获取,不硬编码密码
* 3. 数据验证:过滤无效价格和空值
* 4. 添加降级策略:MySQL不可用时暂存到HDFS
*/
object SalesDataToMySQL {
private val logger: Logger = LoggerFactory.getLogger(SalesDataToMySQL.getClass)
def main(args: Array[String]): Unit = {
try {
val spark = SparkSession.builder()
.appName("SalesDataToMySQL")
.master(ConfigManager.getSparkMaster)
.config("spark.sql.shuffle.partitions", ConfigManager.getSparkShufflePartitions)
.getOrCreate()
import spark.implicits._
// 从Kafka读取销售数据
val salesStream = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", ConfigManager.getKafkaBootstrapServers)
.option("subscribe", ConfigManager.getKafkaTopic("sales-data"))
.option("startingOffsets", "latest")
.option("maxOffsetsPerTrigger", "500")
.load()
// 定义Schema
val schema = new StructType()
.add("sale_id", LongType)
.add("vehicle_id", IntegerType)
.add("customer_id", IntegerType)
.add("dealer_id", IntegerType)
.add("sale_date", StringType)
.add("sale_price", DoubleType)
.add("payment_method", StringType)
.add("region", StringType)
// 解析JSON数据 + 数据验证
val salesData = salesStream
.selectExpr("CAST(value AS STRING)")
.select(from_json($"value", schema).as("data"))
.select("data.*")
.filter($"data".isNotNull)
.filter($"sale_id".isNotNull && $"vehicle_id".isNotNull && $"customer_id".isNotNull)
.filter($"sale_date".isNotNull && $"sale_price".isNotNull)
.filter($"payment_method".isNotNull && $"region".isNotNull)
.filter($"sale_price" > 0) // 价格必须为正数
.withColumn("sale_date", to_date($"sale_date"))
// 写入MySQL(带降级策略)
val query = salesData.writeStream
.foreachBatch { (batchDF: org.apache.spark.sql.DataFrame, batchId: Long) =>
logger.info(s"处理批次: $batchId, 记录数: ${batchDF.count()}")
try {
// 使用ConfigManager获取MySQL配置,不硬编码密码
val tableName = ConfigManager.getDatabaseTable("sales")
val mysqlOptions = ConfigManager.getMySQLOptions(tableName)
batchDF.write
.format("jdbc")
.options(mysqlOptions)
.mode("append")
.save()
logger.info(s"批次 $batchId 成功写入MySQL")
} catch {
case e: Exception =>
logger.error(s"批次 $batchId 写入MySQL失败,启动降级策略", e)
fallbackToHDFS(batchDF, batchId)
}
}
.outputMode(OutputMode.Append())
.option("checkpointLocation", ConfigManager.getSparkCheckpointDir + "/sales-to-mysql")
.start()
query.awaitTermination()
} catch {
case e: Exception =>
logger.error("SalesDataToMySQL运行异常", e)
throw e
} finally {
val s = SparkSession.getActiveSession.orNull
if (s != null) {
s.stop()
logger.info("SparkSession已关闭")
}
}
}
/**
* 降级策略:将数据暂存到HDFS
*/
private def fallbackToHDFS(batchDF: org.apache.spark.sql.DataFrame, batchId: Long): Unit = {
if (!ConfigManager.isDegradationEnabled) {
logger.warn("降级策略未启用,数据可能丢失")
return
}
val fallbackPath = s"${ConfigManager.getFallbackDataPath}/sales_data/batch_$batchId"
logger.info(s"降级策略:将批次 $batchId 数据暂存到 $fallbackPath")
try {
batchDF.write
.mode("overwrite")
.json(fallbackPath)
logger.info(s"批次 $batchId 数据已暂存到HDFS: $fallbackPath")
} catch {
case e: Exception =>
logger.error(s"降级策略执行失败,批次 $batchId 数据可能丢失", e)
}
}
}3. 数据分析模块#
3.1 销量分析#
package com.sparklearning.auto.analysis
import com.sparklearning.auto.config.ConfigManager
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.slf4j.{Logger, LoggerFactory}
/**
* 销量分析 - 基于Spark SQL的离线批处理分析
*
* 分析内容:
* 1. 按地区销量分析
* 2. 按支付方式销量分析
* 3. 销售趋势分析(同比/环比)
* 4. 品牌市场份额分析
*
* 企业级改进:
* 1. try-catch-finally确保SparkSession释放
* 2. MySQL配置从ConfigManager获取
* 3. 数据验证和空值处理
* 4. 完善的异常处理和日志
*/
object SalesAnalysis {
private val logger: Logger = LoggerFactory.getLogger(SalesAnalysis.getClass)
def main(args: Array[String]): Unit = {
try {
val spark = SparkSession.builder()
.appName("SalesAnalysis")
.master(ConfigManager.getSparkMaster)
.config("spark.sql.shuffle.partitions", ConfigManager.getSparkShufflePartitions)
.config("spark.sql.adaptive.enabled", "true")
.enableHiveSupport()
.getOrCreate()
import spark.implicits._
// 从Hive读取销售数据
val salesDF = spark.sql("SELECT * FROM auto_database.sales")
// 数据验证:过滤无效记录
val validSalesDF = salesDF
.filter($"sale_price".isNotNull && $"sale_price" > 0)
.filter($"region".isNotNull && $"region" =!= "")
.filter($"sale_date".isNotNull)
.filter($"payment_method".isNotNull && $"payment_method" =!= "")
logger.info(s"有效销售记录数: ${validSalesDF.count()}")
// ========================================
// 1. 按地区销量分析
// ========================================
val regionSales = validSalesDF
.groupBy($"region")
.agg(
sum($"sale_price").as("total_sales"),
count("*").as("sale_count"),
avg($"sale_price").as("avg_price"),
min($"sale_price").as("min_price"),
max($"sale_price").as("max_price")
)
.orderBy($"total_sales".desc)
println("=== 按地区销售分析 ===")
regionSales.show(false)
// ========================================
// 2. 按支付方式销量分析
// ========================================
val paymentSales = validSalesDF
.groupBy($"payment_method")
.agg(
sum($"sale_price").as("total_sales"),
count("*").as("sale_count"),
avg($"sale_price").as("avg_price")
)
.withColumn("ratio", round($"sale_count" / sum("sale_count").over(), 4))
.orderBy($"total_sales".desc)
println("=== 按支付方式销售分析 ===")
paymentSales.show(false)
// ========================================
// 3. 销售趋势分析(月度同比/环比)
// ========================================
val salesTrend = validSalesDF
.withColumn("year", year($"sale_date"))
.withColumn("month", month($"sale_date"))
.groupBy($"year", $"month")
.agg(
sum($"sale_price").as("total_sales"),
count("*").as("sale_count")
)
.orderBy($"year", $"month")
// 计算环比增长率
val salesTrendWithMoM = salesTrend
.withColumn("prev_month_sales", lag($"total_sales", 1).over(
org.apache.spark.sql.expressions.Window.orderBy($"year", $"month")
))
.withColumn("mom_growth_rate",
when($"prev_month_sales".isNotNull && $"prev_month_sales" =!= 0,
round(($"total_sales" - $"prev_month_sales") / $"prev_month_sales" * 100, 2)
).otherwise(lit(null))
)
.drop("prev_month_sales")
println("=== 销售趋势分析(含环比增长率) ===")
salesTrendWithMoM.show(50, false)
// ========================================
// 4. 品牌市场份额分析
// ========================================
val brandMarketShare = validSalesDF
.groupBy($"brand")
.agg(
sum($"sale_price").as("total_sales"),
count("*").as("sale_count")
)
.withColumn("market_share",
round($"sale_count" / sum("sale_count").over() * 100, 2)
)
.withColumn("revenue_share",
round($"total_sales" / sum("total_sales").over() * 100, 2)
)
.orderBy($"sale_count".desc)
println("=== 品牌市场份额分析 ===")
brandMarketShare.show(false)
// ========================================
// 保存分析结果
// ========================================
saveToMySQL(regionSales, "region-sales-analysis")
saveToMySQL(paymentSales, "payment-sales-analysis")
saveToMySQL(salesTrendWithMoM, "sales-trend-analysis")
saveToMySQL(brandMarketShare, "brand-market-share")
// 同时保存到Hive(批处理层结果)
regionSales.write.mode("overwrite").saveAsTable("auto_database.region_sales_analysis")
paymentSales.write.mode("overwrite").saveAsTable("auto_database.payment_sales_analysis")
salesTrendWithMoM.write.mode("overwrite").saveAsTable("auto_database.sales_trend_analysis")
brandMarketShare.write.mode("overwrite").saveAsTable("auto_database.brand_market_share")
logger.info("销量分析完成,结果已保存到MySQL和Hive")
} catch {
case e: Exception =>
logger.error("销量分析运行异常", e)
throw e
} finally {
val s = SparkSession.getActiveSession.orNull
if (s != null) {
s.stop()
logger.info("SparkSession已关闭")
}
}
}
/**
* 保存DataFrame到MySQL
*/
private def saveToMySQL(df: org.apache.spark.sql.DataFrame, tableKey: String): Unit = {
try {
val tableName = ConfigManager.getDatabaseTable(tableKey)
val mysqlOptions = ConfigManager.getMySQLOptions(tableName)
df.write
.format("jdbc")
.options(mysqlOptions)
.mode("overwrite")
.save()
logger.info(s"分析结果已保存到MySQL表: $tableName")
} catch {
case e: Exception =>
logger.error(s"保存到MySQL失败: $tableKey", e)
// 降级:保存到HDFS
val fallbackPath = s"${ConfigManager.getFallbackDataPath}/analysis/$tableKey"
df.write.mode("overwrite").parquet(fallbackPath)
logger.info(s"降级:分析结果已保存到HDFS: $fallbackPath")
}
}
}3.2 价格分析#
package com.sparklearning.auto.analysis
import com.sparklearning.auto.config.ConfigManager
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.slf4j.{Logger, LoggerFactory}
/**
* 价格分析 - 汽车价格维度深度分析
*
* 分析内容:
* 1. 各品牌平均成交价与中位数价格
* 2. 价格离散度分析
* 3. 优惠幅度分析
* 4. 价格弹性分析
*/
object PriceAnalysis {
private val logger: Logger = LoggerFactory.getLogger(PriceAnalysis.getClass)
def main(args: Array[String]): Unit = {
try {
val spark = SparkSession.builder()
.appName("PriceAnalysis")
.master(ConfigManager.getSparkMaster)
.config("spark.sql.shuffle.partitions", ConfigManager.getSparkShufflePartitions)
.config("spark.sql.adaptive.enabled", "true")
.enableHiveSupport()
.getOrCreate()
import spark.implicits._
// 读取销售和车辆数据
val salesDF = spark.sql("SELECT * FROM auto_database.sales")
.filter($"sale_price".isNotNull && $"sale_price" > 0)
val vehiclesDF = spark.sql("SELECT * FROM auto_database.vehicles")
.filter($"price".isNotNull && $"price" > 0)
// 关联销售和车辆数据
val salesWithVehicle = salesDF
.join(vehiclesDF, Seq("vehicle_id"), "left")
.filter($"brand".isNotNull)
// ========================================
// 1. 各品牌价格水平分析
// ========================================
val brandPriceLevel = salesWithVehicle
.groupBy($"brand")
.agg(
avg($"sale_price").as("avg_sale_price"),
expr("percentile_approx(sale_price, 0.5)").as("median_price"),
min($"sale_price").as("min_price"),
max($"sale_price").as("max_price"),
count("*").as("sale_count")
)
.orderBy($"avg_sale_price".desc)
println("=== 各品牌价格水平分析 ===")
brandPriceLevel.show(false)
// ========================================
// 2. 价格离散度分析
// ========================================
val brandPriceDispersion = salesWithVehicle
.groupBy($"brand")
.agg(
avg($"sale_price").as("avg_price"),
stddev($"sale_price").as("stddev_price"),
(stddev($"sale_price") / avg($"sale_price") * 100).as("cv_percent")
)
.orderBy($"cv_percent".desc)
println("=== 价格离散度分析(变异系数CV) ===")
brandPriceDispersion.show(false)
// ========================================
// 3. 优惠幅度分析
// ========================================
val discountAnalysis = salesWithVehicle
.filter($"price" > 0)
.withColumn("discount_rate",
round(($"price" - $"sale_price") / $"price" * 100, 2)
)
.groupBy($"brand")
.agg(
avg($"discount_rate").as("avg_discount_rate"),
max($"discount_rate").as("max_discount_rate"),
min($"discount_rate").as("min_discount_rate")
)
.orderBy($"avg_discount_rate".desc)
println("=== 优惠幅度分析 ===")
discountAnalysis.show(false)
// ========================================
// 4. 地区价格差异分析
// ========================================
val regionPriceDiff = salesWithVehicle
.groupBy($"region", $"brand")
.agg(
avg($"sale_price").as("avg_price"),
count("*").as("sale_count")
)
.filter($"sale_count" >= 5) // 过滤样本量过小的数据
.orderBy($"brand", $"avg_price".desc)
println("=== 地区价格差异分析 ===")
regionPriceDiff.show(50, false)
// 保存结果
saveToMySQL(brandPriceLevel, "brand-price-level")
saveToMySQL(brandPriceDispersion, "brand-price-dispersion")
saveToMySQL(discountAnalysis, "discount-analysis")
} catch {
case e: Exception =>
logger.error("价格分析运行异常", e)
throw e
} finally {
val s = SparkSession.getActiveSession.orNull
if (s != null) {
s.stop()
logger.info("SparkSession已关闭")
}
}
}
private def saveToMySQL(df: org.apache.spark.sql.DataFrame, tableKey: String): Unit = {
try {
val tableName = ConfigManager.getDatabaseTable(tableKey)
val mysqlOptions = ConfigManager.getMySQLOptions(tableName)
df.write.format("jdbc").options(mysqlOptions).mode("overwrite").save()
logger.info(s"分析结果已保存到MySQL表: $tableName")
} catch {
case e: Exception =>
logger.error(s"保存到MySQL失败: $tableKey", e)
val fallbackPath = s"${ConfigManager.getFallbackDataPath}/analysis/$tableKey"
df.write.mode("overwrite").parquet(fallbackPath)
logger.info(s"降级:分析结果已保存到HDFS: $fallbackPath")
}
}
}3.3 用户画像分析#
package com.sparklearning.auto.analysis
import com.sparklearning.auto.config.ConfigManager
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.slf4j.{Logger, LoggerFactory}
/**
* 用户画像分析 - 客户维度深度分析
*
* 分析内容:
* 1. 客户年龄分布
* 2. 客户收入与购车偏好
* 3. 支付方式偏好分析
* 4. 地区客户特征分析
*/
object CustomerPortraitAnalysis {
private val logger: Logger = LoggerFactory.getLogger(CustomerPortraitAnalysis.getClass)
def main(args: Array[String]): Unit = {
try {
val spark = SparkSession.builder()
.appName("CustomerPortraitAnalysis")
.master(ConfigManager.getSparkMaster)
.config("spark.sql.shuffle.partitions", ConfigManager.getSparkShufflePartitions)
.config("spark.sql.adaptive.enabled", "true")
.enableHiveSupport()
.getOrCreate()
import spark.implicits._
// 读取数据
val salesDF = spark.sql("SELECT * FROM auto_database.sales")
.filter($"sale_price".isNotNull && $"sale_price" > 0)
val customersDF = spark.sql("SELECT * FROM auto_database.customers")
.filter($"age".isNotNull && $"age" >= 18 && $"age" <= 100)
.filter($"gender".isNotNull)
val vehiclesDF = spark.sql("SELECT * FROM auto_database.vehicles")
.filter($"brand".isNotNull)
// 关联数据
val salesWithCustomer = salesDF
.join(customersDF, Seq("customer_id"), "left")
.filter($"name".isNotNull)
val fullDF = salesWithCustomer
.join(vehiclesDF, Seq("vehicle_id"), "left")
// ========================================
// 1. 客户年龄分布
// ========================================
val ageDistribution = fullDF
.withColumn("age_group",
when($"age" < 25, "25岁以下")
.when($"age" < 35, "25-34岁")
.when($"age" < 45, "35-44岁")
.when($"age" < 55, "45-54岁")
.otherwise("55岁以上")
)
.groupBy($"age_group")
.agg(
count("*").as("customer_count"),
avg($"sale_price").as("avg_purchase_price"),
round($"customer_count" / sum("customer_count").over() * 100, 2).as("percentage")
)
.orderBy($"customer_count".desc)
println("=== 客户年龄分布 ===")
ageDistribution.show(false)
// ========================================
// 2. 客户收入与购车偏好
// ========================================
val incomePreference = fullDF
.filter($"income".isNotNull && $"income" > 0)
.withColumn("income_level",
when($"income" < 100000, "10万以下")
.when($"income" < 300000, "10-30万")
.when($"income" < 500000, "30-50万")
.when($"income" < 1000000, "50-100万")
.otherwise("100万以上")
)
.groupBy($"income_level", $"brand")
.agg(count("*").as("purchase_count"))
.withColumn("rank", row_number().over(
org.apache.spark.sql.expressions.Window
.partitionBy($"income_level")
.orderBy($"purchase_count".desc)
))
.filter($"rank" <= 3)
.orderBy($"income_level", $"rank")
println("=== 各收入水平Top3偏好品牌 ===")
incomePreference.show(false)
// ========================================
// 3. 性别与购车偏好
// ========================================
val genderPreference = fullDF
.groupBy($"gender", $"brand")
.agg(
count("*").as("purchase_count"),
avg($"sale_price").as("avg_price")
)
.withColumn("rank", row_number().over(
org.apache.spark.sql.expressions.Window
.partitionBy($"gender")
.orderBy($"purchase_count".desc)
))
.filter($"rank" <= 5)
.orderBy($"gender", $"rank")
println("=== 性别与购车偏好Top5 ===")
genderPreference.show(false)
// ========================================
// 4. 地区客户特征分析
// ========================================
val regionCustomerProfile = fullDF
.groupBy($"region")
.agg(
count("*").as("total_sales"),
avg($"sale_price").as("avg_price"),
avg($"age").as("avg_age"),
avg($"income").as("avg_income"),
round(sum(when($"payment_method" === "贷款", 1).otherwise(0)) / count("*") * 100, 2).as("loan_rate")
)
.orderBy($"total_sales".desc)
println("=== 地区客户特征分析 ===")
regionCustomerProfile.show(false)
// 保存结果
saveToMySQL(ageDistribution, "customer-age-distribution")
saveToMySQL(incomePreference, "income-brand-preference")
saveToMySQL(regionCustomerProfile, "region-customer-profile")
} catch {
case e: Exception =>
logger.error("用户画像分析运行异常", e)
throw e
} finally {
val s = SparkSession.getActiveSession.orNull
if (s != null) {
s.stop()
logger.info("SparkSession已关闭")
}
}
}
private def saveToMySQL(df: org.apache.spark.sql.DataFrame, tableKey: String): Unit = {
try {
val tableName = ConfigManager.getDatabaseTable(tableKey)
val mysqlOptions = ConfigManager.getMySQLOptions(tableName)
df.write.format("jdbc").options(mysqlOptions).mode("overwrite").save()
logger.info(s"分析结果已保存到MySQL表: $tableName")
} catch {
case e: Exception =>
logger.error(s"保存到MySQL失败: $tableKey", e)
val fallbackPath = s"${ConfigManager.getFallbackDataPath}/analysis/$tableKey"
df.write.mode("overwrite").parquet(fallbackPath)
logger.info(s"降级:分析结果已保存到HDFS: $fallbackPath")
}
}
}3.4 传感器数据分析#
package com.sparklearning.auto.analysis
import com.sparklearning.auto.config.ConfigManager
import com.sparklearning.auto.dataaccess.DataAccessLayer
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.hadoop.hbase.io.ImmutableBytesWritable
import org.apache.hadoop.hbase.mapreduce.TableInputFormat
import org.apache.spark.rdd.RDD
import org.slf4j.{Logger, LoggerFactory}
/**
* 传感器数据分析 - 车辆传感器数据深度分析
*
* 分析内容:
* 1. 异常传感器统计
* 2. 传感器值趋势分析
* 3. 车辆健康评分
* 4. 传感器异常率分析
*
* 企业级改进:
* 1. HBase配置从ConfigManager获取
* 2. try-finally确保SparkSession释放
* 3. MySQL配置从ConfigManager获取
* 4. 数据验证和空值处理
*/
object SensorAnalysis {
private val logger: Logger = LoggerFactory.getLogger(SensorAnalysis.getClass)
def main(args: Array[String]): Unit = {
try {
val spark = SparkSession.builder()
.appName("SensorAnalysis")
.master(ConfigManager.getSparkMaster)
.config("spark.sql.shuffle.partitions", ConfigManager.getSparkShufflePartitions)
.config("spark.sql.adaptive.enabled", "true")
.getOrCreate()
// 配置HBase(从ConfigManager获取)
val hbaseConfig = ConfigManager.getHBaseConfig
hbaseConfig.set(TableInputFormat.INPUT_TABLE, ConfigManager.getHBaseTableName("sensor-data"))
// 从HBase读取数据
val hbaseRDD: RDD[(ImmutableBytesWritable, org.apache.hadoop.hbase.client.Result)] =
spark.sparkContext.newAPIHadoopRDD(
hbaseConfig,
classOf[TableInputFormat],
classOf[ImmutableBytesWritable],
classOf[org.apache.hadoop.hbase.client.Result]
)
// 转换为DataFrame(添加空值处理)
val sensorDF = hbaseRDD.flatMap { case (_, result) =>
try {
val rowKey = new String(result.getRow)
val vehicleIdBytes = result.getValue("cf".getBytes, "vehicle_id".getBytes)
val sensorTypeBytes = result.getValue("cf".getBytes, "sensor_type".getBytes)
val valueBytes = result.getValue("cf".getBytes, "value".getBytes)
val statusBytes = result.getValue("cf".getBytes, "status".getBytes)
// 空值检查
if (vehicleIdBytes == null || sensorTypeBytes == null ||
valueBytes == null || statusBytes == null) {
logger.warn(s"跳过空值记录: rowKey=$rowKey")
None
} else {
val vehicleId = new String(vehicleIdBytes).toInt
val sensorType = new String(sensorTypeBytes)
val value = new String(valueBytes).toDouble
val status = new String(statusBytes)
val timestamp = rowKey.split("-")(0).toLong
// 数据验证
if (DataAccessLayer.validateSensorValue(sensorType, value)) {
Some((vehicleId, sensorType, value, status, timestamp))
} else {
logger.warn(s"跳过无效传感器值: type=$sensorType, value=$value")
None
}
}
} catch {
case e: Exception =>
logger.warn("解析HBase记录失败", e)
None
}
}.toDF("vehicle_id", "sensor_type", "value", "status", "timestamp")
// ========================================
// 1. 异常传感器统计
// ========================================
val abnormalSensors = sensorDF
.filter($"status" === "异常")
.groupBy($"sensor_type")
.agg(
count("*").as("abnormal_count"),
avg($"value").as("avg_abnormal_value")
)
.orderBy($"abnormal_count".desc)
println("=== 异常传感器统计 ===")
abnormalSensors.show(false)
// ========================================
// 2. 传感器值统计(正常vs异常对比)
// ========================================
val sensorStats = sensorDF
.groupBy($"sensor_type", $"status")
.agg(
avg($"value").as("avg_value"),
stddev($"value").as("stddev_value"),
max($"value").as("max_value"),
min($"value").as("min_value"),
count("*").as("count")
)
.orderBy($"sensor_type", $"status")
println("=== 传感器值统计(正常vs异常) ===")
sensorStats.show(false)
// ========================================
// 3. 车辆健康评分
// ========================================
val vehicleHealthScore = sensorDF
.groupBy($"vehicle_id")
.agg(
count("*").as("total_readings"),
sum(when($"status" === "异常", 1).otherwise(0)).as("abnormal_count")
)
.withColumn("abnormal_rate", $"abnormal_count" / $"total_readings")
.withColumn("health_score",
when($"abnormal_rate" < 0.05, 100)
.when($"abnormal_rate" < 0.10, 90)
.when($"abnormal_rate" < 0.20, 80)
.when($"abnormal_rate" < 0.30, 70)
.when($"abnormal_rate" < 0.50, 60)
.otherwise(50)
)
.orderBy($"health_score".asc)
println("=== 车辆健康评分(按评分升序) ===")
vehicleHealthScore.show(20, false)
// ========================================
// 4. 传感器异常率排名
// ========================================
val sensorAbnormalRate = sensorDF
.groupBy($"sensor_type")
.agg(
count("*").as("total_count"),
sum(when($"status" === "异常", 1).otherwise(0)).as("abnormal_count")
)
.withColumn("abnormal_rate",
round($"abnormal_count" / $"total_count" * 100, 2)
)
.orderBy($"abnormal_rate".desc)
println("=== 传感器异常率排名 ===")
sensorAbnormalRate.show(false)
// ========================================
// 保存分析结果到MySQL
// ========================================
saveToMySQL(abnormalSensors, "sensor-abnormal-analysis")
saveToMySQL(sensorStats, "sensor-stats")
saveToMySQL(vehicleHealthScore, "vehicle-health-score")
saveToMySQL(sensorAbnormalRate, "sensor-abnormal-rate")
} catch {
case e: Exception =>
logger.error("传感器分析运行异常", e)
throw e
} finally {
val s = SparkSession.getActiveSession.orNull
if (s != null) {
s.stop()
logger.info("SparkSession已关闭")
}
}
}
private def saveToMySQL(df: org.apache.spark.sql.DataFrame, tableKey: String): Unit = {
try {
val tableName = ConfigManager.getDatabaseTable(tableKey)
val mysqlOptions = ConfigManager.getMySQLOptions(tableName)
df.write.format("jdbc").options(mysqlOptions).mode("overwrite").save()
logger.info(s"分析结果已保存到MySQL表: $tableName")
} catch {
case e: Exception =>
logger.error(s"保存到MySQL失败: $tableKey", e)
val fallbackPath = s"${ConfigManager.getFallbackDataPath}/analysis/$tableKey"
df.write.mode("overwrite").parquet(fallbackPath)
logger.info(s"降级:分析结果已保存到HDFS: $fallbackPath")
}
}
}3.5 竞争力分析#
package com.sparklearning.auto.analysis
import com.sparklearning.auto.config.ConfigManager
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.slf4j.{Logger, LoggerFactory}
/**
* 竞争力分析 - 品牌竞争力维度分析
*
* 分析内容:
* 1. 市场占有率分析
* 2. 品牌排名分析
* 3. 地区竞争力对比
*/
object CompetitivenessAnalysis {
private val logger: Logger = LoggerFactory.getLogger(CompetitivenessAnalysis.getClass)
def main(args: Array[String]): Unit = {
try {
val spark = SparkSession.builder()
.appName("CompetitivenessAnalysis")
.master(ConfigManager.getSparkMaster)
.config("spark.sql.shuffle.partitions", ConfigManager.getSparkShufflePartitions)
.config("spark.sql.adaptive.enabled", "true")
.enableHiveSupport()
.getOrCreate()
import spark.implicits._
val salesDF = spark.sql("SELECT * FROM auto_database.sales")
.filter($"sale_price".isNotNull && $"sale_price" > 0)
val vehiclesDF = spark.sql("SELECT * FROM auto_database.vehicles")
.filter($"brand".isNotNull)
val fullDF = salesDF.join(vehiclesDF, Seq("vehicle_id"), "left")
.filter($"brand".isNotNull)
// ========================================
// 1. 市场占有率分析
// ========================================
val marketShare = fullDF
.groupBy($"brand")
.agg(
count("*").as("sale_count"),
sum($"sale_price").as("total_revenue")
)
.withColumn("volume_share", round($"sale_count" / sum("sale_count").over() * 100, 2))
.withColumn("revenue_share", round($"total_revenue" / sum("total_revenue").over() * 100, 2))
.withColumn("volume_rank", row_number().over(
org.apache.spark.sql.expressions.Window.orderBy($"sale_count".desc)
))
.withColumn("revenue_rank", row_number().over(
org.apache.spark.sql.expressions.Window.orderBy($"total_revenue".desc)
))
.orderBy($"volume_share".desc)
println("=== 市场占有率分析 ===")
marketShare.show(false)
// ========================================
// 2. 地区竞争力对比
// ========================================
val regionCompetitiveness = fullDF
.groupBy($"region", $"brand")
.agg(count("*").as("sale_count"))
.withColumn("market_share",
round($"sale_count" / sum("sale_count").over(
org.apache.spark.sql.expressions.Window.partitionBy($"region")
) * 100, 2)
)
.withColumn("rank", row_number().over(
org.apache.spark.sql.expressions.Window
.partitionBy($"region")
.orderBy($"sale_count".desc)
))
.filter($"rank" <= 3)
.orderBy($"region", $"rank")
println("=== 各地区Top3品牌 ===")
regionCompetitiveness.show(50, false)
// 保存结果
saveToMySQL(marketShare, "vehicle-competitiveness")
} catch {
case e: Exception =>
logger.error("竞争力分析运行异常", e)
throw e
} finally {
val s = SparkSession.getActiveSession.orNull
if (s != null) {
s.stop()
logger.info("SparkSession已关闭")
}
}
}
private def saveToMySQL(df: org.apache.spark.sql.DataFrame, tableKey: String): Unit = {
try {
val tableName = ConfigManager.getDatabaseTable(tableKey)
val mysqlOptions = ConfigManager.getMySQLOptions(tableName)
df.write.format("jdbc").options(mysqlOptions).mode("overwrite").save()
logger.info(s"分析结果已保存到MySQL表: $tableName")
} catch {
case e: Exception =>
logger.error(s"保存到MySQL失败: $tableKey", e)
val fallbackPath = s"${ConfigManager.getFallbackDataPath}/analysis/$tableKey"
df.write.mode("overwrite").parquet(fallbackPath)
logger.info(s"降级:分析结果已保存到HDFS: $fallbackPath")
}
}
}4. 实时流处理模块(速度层)#
4.1 实时传感器告警#
package com.sparklearning.auto.streaming
import com.sparklearning.auto.config.ConfigManager
import com.sparklearning.auto.dataaccess.DataAccessLayer
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.OutputMode
import org.apache.spark.sql.types._
import org.slf4j.{Logger, LoggerFactory}
/**
* 实时传感器告警 - 速度层实时处理
*
* 功能:
* 1. 实时检测传感器异常
* 2. 窗口内异常率统计
* 3. 异常告警发送到Kafka告警主题
* 4. 降级策略:流处理失败时回退到批处理
*
* 知识点应用:
* - DStream微批次处理
* - 窗口计算(5秒批次,30秒窗口,10秒滑动)
* - 背压机制
*/
object RealtimeSensorAlert {
private val logger: Logger = LoggerFactory.getLogger(RealtimeSensorAlert.getClass)
// 连续失败计数器(用于降级判断)
private var consecutiveFailures = 0
def main(args: Array[String]): Unit = {
try {
val spark = SparkSession.builder()
.appName("RealtimeSensorAlert")
.master(ConfigManager.getSparkMaster)
// 启用背压机制
.config("spark.streaming.backpressure.enabled", "true")
.config("spark.streaming.backpressure.initialRate",
ConfigManager.getSparkBatchDurationSeconds.toString)
.config("spark.streaming.kafka.maxRatePerPartition", "100")
.config("spark.sql.shuffle.partitions", ConfigManager.getSparkShufflePartitions)
.getOrCreate()
import spark.implicits._
// 从Kafka读取传感器数据
val sensorStream = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", ConfigManager.getKafkaBootstrapServers)
.option("subscribe", ConfigManager.getKafkaTopic("sensor-data"))
.option("startingOffsets", "latest")
.load()
// 定义Schema
val schema = new StructType()
.add("sensor_id", IntegerType)
.add("vehicle_id", IntegerType)
.add("timestamp", LongType)
.add("sensor_type", StringType)
.add("value", DoubleType)
.add("status", StringType)
// 解析JSON + 数据验证
val parsedStream = sensorStream
.selectExpr("CAST(value AS STRING)", "timestamp as kafka_timestamp")
.select(from_json($"value", schema).as("data"), $"kafka_timestamp")
.select("data.*", "kafka_timestamp")
.filter($"data".isNotNull)
.filter($"sensor_id".isNotNull && $"vehicle_id".isNotNull)
.filter($"sensor_type".isNotNull && $"value".isNotNull)
.withColumn("event_time", to_timestamp($"timestamp" / 1000))
.withWatermark("event_time", "10 seconds")
// ========================================
// 1. 实时异常检测(状态过滤)
// ========================================
val abnormalStream = parsedStream
.filter($"status" === "异常")
// ========================================
// 2. 窗口内异常率统计
// 窗口长度30秒,滑动步长10秒
// ========================================
val windowedAbnormalRate = parsedStream
.groupBy(
window($"event_time", "30 seconds", "10 seconds"),
$"sensor_type"
)
.agg(
count("*").as("total_count"),
sum(when($"status" === "异常", 1).otherwise(0)).as("abnormal_count")
)
.withColumn("abnormal_rate",
round($"abnormal_count" / $"total_count" * 100, 2)
)
.filter($"abnormal_rate" > 10) // 异常率超过10%触发告警
// ========================================
// 3. 车辆级别异常聚合
// ========================================
val vehicleAbnormal = abnormalStream
.groupBy(
window($"event_time", "30 seconds", "10 seconds"),
$"vehicle_id"
)
.agg(
count("*").as("abnormal_count"),
collect_set($"sensor_type").as("abnormal_sensor_types")
)
.filter($"abnormal_count" >= 3) // 单车30秒内3次以上异常触发告警
// ========================================
// 输出1:异常告警写入Kafka告警主题
// ========================================
val alertQuery = windowedAbnormalRate
.select(
to_json(struct(
$"window",
$"sensor_type",
$"total_count",
$"abnormal_count",
$"abnormal_rate"
)).as("value")
)
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", ConfigManager.getKafkaBootstrapServers)
.option("topic", ConfigManager.getKafkaTopic("alert-data"))
.option("checkpointLocation", ConfigManager.getSparkCheckpointDir + "/alert-to-kafka")
.outputMode(OutputMode.Update())
.start()
// ========================================
// 输出2:异常率统计写入MySQL
// ========================================
val mysqlQuery = windowedAbnormalRate.writeStream
.foreachBatch { (batchDF: org.apache.spark.sql.DataFrame, batchId: Long) =>
try {
val tableName = ConfigManager.getDatabaseTable("sensor-abnormal-analysis")
val mysqlOptions = ConfigManager.getMySQLOptions(tableName)
batchDF.write
.format("jdbc")
.options(mysqlOptions)
.mode("append")
.save()
consecutiveFailures = 0 // 重置失败计数
logger.info(s"实时告警批次 $batchId 写入MySQL成功")
} catch {
case e: Exception =>
consecutiveFailures += 1
logger.error(s"实时告警批次 $batchId 写入MySQL失败(连续失败: $consecutiveFailures)", e)
// 降级判断
if (consecutiveFailures >= ConfigManager.getBatchFallbackThreshold) {
logger.warn(s"连续失败次数达到阈值(${ConfigManager.getBatchFallbackThreshold}),建议切换到批处理模式")
}
// 降级:暂存到HDFS
val fallbackPath = s"${ConfigManager.getFallbackDataPath}/realtime-alert/batch_$batchId"
batchDF.write.mode("overwrite").json(fallbackPath)
}
}
.outputMode(OutputMode.Update())
.option("checkpointLocation", ConfigManager.getSparkCheckpointDir + "/alert-to-mysql")
.start()
// ========================================
// 输出3:控制台输出(调试用)
// ========================================
val consoleQuery = vehicleAbnormal.writeStream
.format("console")
.option("truncate", "false")
.outputMode(OutputMode.Update())
.start()
spark.streams.awaitAnyTermination()
} catch {
case e: Exception =>
logger.error("实时传感器告警运行异常", e)
throw e
} finally {
val s = SparkSession.getActiveSession.orNull
if (s != null) {
s.stop()
logger.info("SparkSession已关闭")
}
}
}
}5. 降级策略模块#
package com.sparklearning.auto.degradation
import com.sparklearning.auto.config.ConfigManager
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.slf4j.{Logger, LoggerFactory}
/**
* 降级策略管理器
*
* 降级策略设计:
* 1. 实时处理连续失败超过阈值时,自动切换到批处理模式
* 2. 降级期间数据暂存到HDFS,等待恢复后补算
* 3. 提供手动恢复和自动恢复两种模式
*
* 降级流程:
* 实时处理正常 --> 实时处理异常 --> 重试(N次) --> 降级到批处理 --> 恢复检测 --> 恢复实时处理
*/
object DegradationManager {
private val logger: Logger = LoggerFactory.getLogger(DegradationManager.getClass)
/**
* 执行降级批处理 - 处理暂存到HDFS的降级数据
*
* 当实时处理恢复正常后,调用此方法补算降级期间的数据
*/
def processFallbackData(): Unit = {
if (!ConfigManager.isDegradationEnabled) {
logger.info("降级策略未启用,跳过降级数据处理")
return
}
try {
val spark = SparkSession.builder()
.appName("DegradationBatchProcessor")
.master(ConfigManager.getSparkMaster)
.config("spark.sql.shuffle.partitions", ConfigManager.getSparkShufflePartitions)
.getOrCreate()
val fallbackBasePath = ConfigManager.getFallbackDataPath
// 处理传感器降级数据
processSensorFallbackData(spark, fallbackBasePath)
// 处理销售降级数据
processSalesFallbackData(spark, fallbackBasePath)
logger.info("降级数据处理完成")
} catch {
case e: Exception =>
logger.error("降级数据处理异常", e)
throw e
} finally {
val s = SparkSession.getActiveSession.orNull
if (s != null) {
s.stop()
}
}
}
/**
* 处理传感器降级数据
*/
private def processSensorFallbackData(spark: SparkSession, basePath: String): Unit = {
val sensorFallbackPath = s"$basePath/sensor_data"
try {
val sensorDF = spark.read.json(sensorFallbackPath)
if (sensorDF.count() > 0) {
logger.info(s"发现传感器降级数据: ${sensorDF.count()} 条")
// 写入HBase
val tableName = ConfigManager.getHBaseTableName("sensor-data")
val hbaseConfig = ConfigManager.getHBaseConfig
hbaseConfig.set("hbase.mapreduce.inputtable", tableName)
// 使用DataAccessLayer写入
val records = sensorDF.collect().flatMap { row =>
try {
val sensorType = row.getAs[String]("sensor_type")
val value = row.getAs[Double]("value")
if (sensorType != null && value != null) {
val rowKey = s"${row.getAs[Long]("timestamp")}-${row.getAs[Int]("sensor_id")}"
val columns = Seq(
("vehicle_id", row.getAs[Int]("vehicle_id").asInstanceOf[Any]),
("sensor_type", sensorType.asInstanceOf[Any]),
("value", value.asInstanceOf[Any]),
("status", row.getAs[String]("status").asInstanceOf[Any])
)
Some((rowKey, columns))
} else None
} catch {
case _: Exception => None
}
}
if (records.nonEmpty) {
val result = com.sparklearning.auto.dataaccess.DataAccessLayer.batchInsertToHBase(
"sensor-data", "cf", records
)
result match {
case scala.util.Success(count) =>
logger.info(s"降级传感器数据已写入HBase: $count 条")
case scala.util.Failure(e) =>
logger.error("降级传感器数据写入HBase失败", e)
}
}
}
} catch {
case e: Exception =>
logger.warn("未找到传感器降级数据或处理失败", e)
}
}
/**
* 处理销售降级数据
*/
private def processSalesFallbackData(spark: SparkSession, basePath: String): Unit = {
val salesFallbackPath = s"$basePath/sales_data"
try {
val salesDF = spark.read.json(salesFallbackPath)
if (salesDF.count() > 0) {
logger.info(s"发现销售降级数据: ${salesDF.count()} 条")
val tableName = ConfigManager.getDatabaseTable("sales")
val mysqlOptions = ConfigManager.getMySQLOptions(tableName)
salesDF.write
.format("jdbc")
.options(mysqlOptions)
.mode("append")
.save()
logger.info("降级销售数据已写入MySQL")
}
} catch {
case e: Exception =>
logger.warn("未找到销售降级数据或处理失败", e)
}
}
/**
* 带重试的操作执行器
*
* @param operation 要执行的操作
* @param maxRetries 最大重试次数
* @param retryIntervalMs 重试间隔(毫秒)
* @tparam T 返回值类型
* @return 操作结果
*/
def withRetry[T](operation: => T,
maxRetries: Int = ConfigManager.getMaxRetries,
retryIntervalMs: Long = ConfigManager.getRetryIntervalMs): Try[T] = {
var lastException: Option[Exception] = None
var attempt = 0
while (attempt < maxRetries) {
try {
return scala.util.Success(operation)
} catch {
case e: Exception =>
lastException = Some(e)
attempt += 1
logger.warn(s"操作失败,第 $attempt 次重试(共 $maxRetries 次)", e)
if (attempt < maxRetries) {
Thread.sleep(retryIntervalMs)
}
}
}
scala.util.Failure(lastException.getOrElse(
new RuntimeException(s"操作在 $maxRetries 次重试后仍失败")
))
}
}6. 可视化模块#
6.1 Grafana仪表板#
配置Grafana连接MySQL和HBase数据源,创建以下仪表板:
- 销售概览仪表板
- 总销售额、销售数量、平均售价
- 按地区销售分布
- 销售趋势图表
- 支付方式分布
- 传感器监控仪表板
- 异常传感器数量
- 传感器值实时监控
- 车辆状态分布
- 传感器类型分布
- 用户画像仪表板
- 客户年龄分布
- 收入与购车偏好
- 地区客户特征
- 竞争力仪表板
- 品牌市场份额
- 地区品牌排名
- 价格竞争力对比
6.2 数据可视化代码#
# dashboard.py
# 注意:数据库密码通过环境变量获取,不硬编码
import os
import pandas as pd
import matplotlib.pyplot as plt
import seaborn as sns
import mysql.connector
# 从环境变量获取数据库配置
MYSQL_HOST = os.environ.get('MYSQL_HOST', 'localhost')
MYSQL_USER = os.environ.get('MYSQL_USER', 'auto_admin')
MYSQL_PASSWORD = os.environ.get('MYSQL_PASSWORD', '')
MYSQL_DATABASE = os.environ.get('MYSQL_DATABASE', 'auto_database')
# 连接MySQL(密码从环境变量获取)
conn = mysql.connector.connect(
host=MYSQL_HOST,
user=MYSQL_USER,
password=MYSQL_PASSWORD,
database=MYSQL_DATABASE
)
try:
# 读取销售数据
sales_df = pd.read_sql("SELECT * FROM sales", conn)
# 读取传感器异常数据
sensor_df = pd.read_sql("SELECT * FROM sensor_abnormal_analysis", conn)
# 设置中文字体
plt.rcParams['font.sans-serif'] = ['SimHei']
plt.rcParams['axes.unicode_minus'] = False
# 1. 销售趋势图
fig, axes = plt.subplots(2, 1, figsize=(12, 8))
sales_df['sale_date'] = pd.to_datetime(sales_df['sale_date'])
sales_by_month = sales_df.groupby(sales_df['sale_date'].dt.to_period('M')).agg({
'sale_price': 'sum',
'sale_id': 'count'
})
sales_by_month['sale_price'].plot(kind='bar', ax=axes[0], title='月度销售额')
sales_by_month['sale_id'].plot(kind='bar', ax=axes[1], title='月度销量')
plt.tight_layout()
plt.savefig('sales_trend.png')
# 2. 地区销售分布图
plt.figure(figsize=(12, 6))
region_sales = sales_df.groupby('region')['sale_price'].sum().sort_values(ascending=False)
region_sales.plot(kind='pie', autopct='%1.1f%%')
plt.title('地区销售分布')
plt.savefig('region_sales.png')
# 3. 传感器异常分布图
plt.figure(figsize=(12, 6))
sns.barplot(x='sensor_type', y='abnormal_count', data=sensor_df)
plt.xticks(rotation=45)
plt.title('传感器异常分布')
plt.tight_layout()
plt.savefig('sensor_abnormal.png')
finally:
conn.close()系统部署与监控#
1. 环境搭建#
1.1 基础环境要求#
| 组件 | 版本 | 最低配置 | 推荐配置 |
|---|---|---|---|
| JDK | 11 | - | OpenJDK 11.0.20+ |
| Scala | 2.13.8 | - | SDKMAN安装 |
| Hadoop | 3.3.6 | 3节点, 8GB/节点 | 5节点, 16GB/节点 |
| Spark | 3.5.8 | 3节点, 8GB/节点 | 5节点, 16GB/节点 |
| Kafka | 3.7.0 | 3 Broker | 5 Broker |
| HBase | 2.5.9 | 3节点 | 5节点 |
| MySQL | 8.0 | 4GB | 8GB |
1.2 环境变量配置#
部署前需设置以下环境变量:
# MySQL配置
export MYSQL_URL="jdbc:mysql://your-mysql-host:3306/auto_database?useSSL=true&serverTimezone=Asia/Shanghai"
export MYSQL_USER="auto_admin"
export MYSQL_PASSWORD="your_secure_password"
# Kafka配置
export KAFKA_BOOTSTRAP_SERVERS="kafka1:9092,kafka2:9092,kafka3:9092"
# HBase配置
export HBASE_ZOOKEEPER_QUORUM="zk1,zk2,zk3"
export HBASE_ZOOKEEPER_CLIENT_PORT="2181"
# Spark配置
export SPARK_MASTER="spark://spark-master:7077"
# 告警配置
export ALERT_EMAIL="admin@your-company.com"1.3 数据库初始化#
-- 创建数据库
CREATE DATABASE IF NOT EXISTS auto_database
DEFAULT CHARACTER SET utf8mb4
DEFAULT COLLATE utf8mb4_unicode_ci;
USE auto_database;
-- 销售数据表
CREATE TABLE IF NOT EXISTS sales (
sale_id BIGINT PRIMARY KEY,
vehicle_id INT NOT NULL,
customer_id INT NOT NULL,
dealer_id INT NOT NULL,
sale_date DATE NOT NULL,
sale_price DECIMAL(12,2) NOT NULL CHECK (sale_price > 0),
payment_method VARCHAR(50) NOT NULL,
region VARCHAR(100) NOT NULL,
create_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
INDEX idx_sale_date (sale_date),
INDEX idx_region (region),
INDEX idx_vehicle_id (vehicle_id)
) ENGINE=InnoDB;
-- 车辆数据表
CREATE TABLE IF NOT EXISTS vehicles (
vehicle_id INT PRIMARY KEY,
brand VARCHAR(50) NOT NULL,
model VARCHAR(100) NOT NULL,
year INT CHECK (year >= 2000 AND year <= 2030),
fuel_type VARCHAR(20) NOT NULL,
engine_size DECIMAL(3,1) CHECK (engine_size > 0),
horsepower INT CHECK (horsepower > 0),
price DECIMAL(12,2) CHECK (price > 0),
features JSON,
INDEX idx_brand (brand)
) ENGINE=InnoDB;
-- 客户数据表
CREATE TABLE IF NOT EXISTS customers (
customer_id INT PRIMARY KEY,
name VARCHAR(100) NOT NULL,
age INT CHECK (age >= 18 AND age <= 100),
gender VARCHAR(10) NOT NULL,
occupation VARCHAR(100),
income DECIMAL(12,2) CHECK (income >= 0),
region VARCHAR(100) NOT NULL,
purchase_history JSON,
INDEX idx_age (age),
INDEX idx_region (region)
) ENGINE=InnoDB;
-- 分析结果表
CREATE TABLE IF NOT EXISTS region_sales_analysis (
region VARCHAR(100),
total_sales DECIMAL(15,2),
sale_count BIGINT,
avg_price DECIMAL(12,2),
min_price DECIMAL(12,2),
max_price DECIMAL(12,2)
) ENGINE=InnoDB;
CREATE TABLE IF NOT EXISTS sensor_abnormal_analysis (
sensor_type VARCHAR(50),
abnormal_count BIGINT,
avg_abnormal_value DOUBLE
) ENGINE=InnoDB;
CREATE TABLE IF NOT EXISTS brand_market_share (
brand VARCHAR(50),
sale_count BIGINT,
total_revenue DECIMAL(15,2),
volume_share DOUBLE,
revenue_share DOUBLE,
volume_rank INT,
revenue_rank INT
) ENGINE=InnoDB;2. Kafka主题创建#
# 创建销售数据主题(3分区,2副本)
kafka-topics.sh --create \
--bootstrap-server localhost:9092 \
--topic sales-data \
--partitions 3 \
--replication-factor 2
# 创建传感器数据主题(6分区,2副本,保留7天数据)
kafka-topics.sh --create \
--bootstrap-server localhost:9092 \
--topic sensor-data \
--partitions 6 \
--replication-factor 2 \
--config retention.ms=604800000
# 创建告警数据主题(3分区,2副本)
kafka-topics.sh --create \
--bootstrap-server localhost:9092 \
--topic alert-data \
--partitions 3 \
--replication-factor 23. HBase表创建#
# 进入HBase Shell
hbase shell
# 创建传感器数据表
create 'sensor_data', {NAME => 'cf', VERSIONS => 1, COMPRESSION => 'SNAPPY'}, {SPLITS => ['1', '2', '3', '4', '5', '6', '7', '8', '9']}
# 创建GPS轨迹表
create 'gps_track', {NAME => 'cf', VERSIONS => 1, COMPRESSION => 'SNAPPY'}4. 系统监控#
4.1 Prometheus监控配置#
# prometheus.yml
scrape_configs:
- job_name: 'spark'
static_configs:
- targets: ['spark-master:4040', 'spark-worker1:4041']
- job_name: 'kafka'
static_configs:
- targets: ['kafka1:9308', 'kafka2:9308', 'kafka3:9308']
- job_name: 'hbase'
static_configs:
- targets: ['hbase-master:16010']4.2 关键监控指标#
| 监控对象 | 关键指标 | 告警阈值 |
|---|---|---|
| Spark作业 | 任务执行时间 | > 5分钟 |
| Spark作业 | 任务失败率 | > 5% |
| Kafka | 消息积压(Lag) | > 10000条 |
| Kafka | 消费延迟 | > 60秒 |
| HBase | 读写延迟 | > 100ms |
| MySQL | 连接数 | > 80%最大连接数 |
| 系统 | CPU使用率 | > 80% |
| 系统 | 内存使用率 | > 85% |
| 系统 | 磁盘使用率 | > 90% |
项目实现步骤#
1. 环境搭建#
- 安装和配置JDK 11、Scala 2.13.8
- 部署Hadoop 3.3.6集群
- 部署Spark 3.5.8集群
- 部署Kafka 3.7.0集群
- 部署HBase 2.5.9集群
- 安装MySQL 8.0并执行初始化SQL
- 配置环境变量(数据库密码等敏感信息)
2. 项目构建#
- 使用Maven构建项目:
mvn clean package - 上传jar包到Spark集群
- 验证
application.conf配置正确性
3. 数据采集#
- 启动Kafka服务和创建主题
- 运行
SalesDataProducer生成模拟销售数据 - 运行
SensorDataProducer生成模拟传感器数据 - 使用Kafka Console Consumer验证数据
4. 数据存储#
- 运行
SensorDataToHBase将传感器数据写入HBase - 运行
SalesDataToMySQL将销售数据写入MySQL - 验证数据正确存储
- 测试降级策略(模拟HBase/MySQL不可用)
5. 数据分析#
- 运行
SalesAnalysis进行销量分析 - 运行
PriceAnalysis进行价格分析 - 运行
CustomerPortraitAnalysis进行用户画像分析 - 运行
SensorAnalysis进行传感器分析 - 运行
CompetitivenessAnalysis进行竞争力分析
6. 实时处理#
- 运行
RealtimeSensorAlert启动实时告警 - 验证窗口计算结果
- 测试背压机制
- 测试降级策略
7. 可视化展示#
- 配置Grafana数据源
- 创建仪表板和图表
- 验证可视化效果
8. 系统监控#
- 配置Prometheus和Grafana监控
- 设置告警规则
- 验证监控效果
项目扩展#
- 实时推荐系统:基于客户购买历史和行为数据,使用Spark MLlib实现车辆推荐
- 预测分析:使用时间序列模型预测销售趋势,使用分类模型预测传感器故障
- 智能诊断:基于OBD诊断数据和传感器数据,实现车辆故障智能诊断
- 供应链优化:分析销售数据,使用优化算法优化车辆供应链
- 驾驶行为评分:基于GPS轨迹和加速度数据,构建驾驶行为评分模型
- 动态定价:基于市场供需和竞争分析,实现动态定价策略
项目总结#
本项目实现了一个完整的企业级汽车大数据分析系统,涵盖了数据采集、存储、处理、分析和可视化的全流程。通过整合Spark、Hadoop、HBase、Kafka等大数据技术,系统能够实时处理和分析汽车行业的各类数据,为企业决策提供数据支持。
系统特点#
| 特点 | 实现方式 | 对应知识点 |
|---|---|---|
| 实时性 | Kafka + Spark Streaming微批次处理 | 知识点3:微批次处理原理 |
| 准确性 | Lambda架构批处理层全量重算 | 知识点2:Lambda架构 |
| 可靠性 | 降级策略 + 数据暂存 + 重试机制 | 降级策略模块 |
| 安全性 | 配置外部化 + 环境变量 + 数据验证 | ConfigManager + DataAccessLayer |
| 可扩展性 | 集群部署 + 水平扩展 + 配置驱动 | application.conf |
| 可观测性 | Prometheus + Grafana + 日志体系 | 监控模块 |
关键改进对比#
| 改进项 | 改进前 | 改进后 |
|---|---|---|
| 配置管理 | 硬编码在代码中 | ConfigManager + application.conf + 环境变量 |
| 密码安全 | password=123456硬编码 | 环境变量注入,配置文件中不存储明文 |
| 资源管理 | 无保障 | try-catch-finally确保SparkSession/Connection释放 |
| 数据验证 | 无 | 空值检查、范围检查、类型检查 |
| 异常处理 | 简单print | SLF4J日志 + 降级策略 + 重试机制 |
| 集合转换 | JavaConverters(已废弃) | scala.jdk.CollectionConverters |
| 降级策略 | 无 | HDFS暂存 + 批处理回退 + 自动恢复 |
| 分析维度 | 仅销量和传感器 | 销量/价格/用户画像/竞争力四维指标体系 |
| 技术栈 | 版本混乱 | Spark 3.5.8 + Scala 2.13.8 + JDK 11统一 |
通过本项目,我们掌握了企业级大数据系统的设计和实现方法,理解了车联网数据特征、Lambda架构、Spark Streaming微批次处理原理以及汽车数据分析指标体系,为未来的大数据项目开发打下了坚实的基础。