跳过正文
  1. Posts/

汽车大数据分析系统

作者
John Lee
Building things with code. Writing about tech, projects, and ideas.
目录

汽车大数据分析系统
#

项目概述
#

汽车大数据分析系统是一个综合性的企业级大数据解决方案,旨在通过收集、存储和分析汽车行业的各类数据,为汽车制造商、经销商和消费者提供有价值的洞察。本系统整合了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架构,原因如下:

  1. 汽车销售数据分析需要精确的历史统计(批处理层保障准确性)
  2. 传感器监控需要实时告警(速度层保障实时性)
  3. 两种数据处理逻辑差异较大,分离实现更清晰
  4. 便于初学者理解批处理与流处理的区别

知识点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-4

DStream(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.enabledfalse是否启用背压
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(品牌销量)品牌/地区
产品竞争力性价比指数配置得分/成交价格品牌/车型
产品竞争力口碑指数评分加权平均品牌/车型
渠道能力经销商覆盖率有销量的经销商数/总经销商数品牌/地区

系统架构
#

技术栈版本
#

组件版本用途
JDK11运行环境
Scala2.13.8开发语言
Spark3.5.8数据处理引擎
Kafka3.7.0消息队列
HBase2.5.9时序数据存储
Hadoop3.3.6分布式存储
Hive3.1.3数据仓库
MySQL8.0业务数据存储
Typesafe Config1.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_idBIGINT销售记录IDPRIMARY KEY
vehicle_idINT车辆IDNOT NULL
customer_idINT客户IDNOT NULL
dealer_idINT经销商IDNOT NULL
sale_dateDATE销售日期NOT NULL
sale_priceDECIMAL(12,2)销售价格CHECK(sale_price > 0)
payment_methodVARCHAR(50)支付方式NOT NULL
regionVARCHAR(100)销售地区NOT NULL

2. 车辆数据模型
#

字段名数据类型描述约束
vehicle_idINT车辆IDPRIMARY KEY
brandVARCHAR(50)品牌NOT NULL
modelVARCHAR(100)型号NOT NULL
yearINT年份CHECK(year >= 2000 AND year <= 2030)
fuel_typeVARCHAR(20)燃油类型NOT NULL
engine_sizeDECIMAL(3,1)发动机排量CHECK(engine_size > 0)
horsepowerINT马力CHECK(horsepower > 0)
priceDECIMAL(12,2)指导价格CHECK(price > 0)
featuresJSON配置特征

3. 客户数据模型
#

字段名数据类型描述约束
customer_idINT客户IDPRIMARY KEY
nameVARCHAR(100)姓名NOT NULL
ageINT年龄CHECK(age >= 18 AND age <= 100)
genderVARCHAR(10)性别NOT NULL
occupationVARCHAR(100)职业
incomeDECIMAL(12,2)收入CHECK(income >= 0)
regionVARCHAR(100)地区NOT NULL
purchase_historyJSON购买历史

4. 传感器数据模型
#

字段名数据类型描述约束
sensor_idINT传感器IDNOT NULL
vehicle_idINT车辆IDNOT NULL
timestampBIGINT采集时间(毫秒时间戳)NOT NULL
sensor_typeVARCHAR(50)传感器类型NOT NULL
valueDOUBLE传感器值NOT NULL
statusVARCHAR(20)状态NOT NULL

5. GPS轨迹数据模型
#

字段名数据类型描述约束
vehicle_idINT车辆IDNOT NULL
timestampBIGINT采集时间(毫秒时间戳)NOT NULL
latitudeDOUBLE纬度CHECK(latitude BETWEEN -90 AND 90)
longitudeDOUBLE经度CHECK(longitude BETWEEN -180 AND 180)
speedDOUBLE速度(km/h)CHECK(speed >= 0)
headingDOUBLE航向角(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数据源,创建以下仪表板:

  1. 销售概览仪表板
  • 总销售额、销售数量、平均售价
  • 按地区销售分布
  • 销售趋势图表
  • 支付方式分布
  1. 传感器监控仪表板
  • 异常传感器数量
  • 传感器值实时监控
  • 车辆状态分布
  • 传感器类型分布
  1. 用户画像仪表板
  • 客户年龄分布
  • 收入与购车偏好
  • 地区客户特征
  1. 竞争力仪表板
  • 品牌市场份额
  • 地区品牌排名
  • 价格竞争力对比

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 基础环境要求
#

组件版本最低配置推荐配置
JDK11-OpenJDK 11.0.20+
Scala2.13.8-SDKMAN安装
Hadoop3.3.63节点, 8GB/节点5节点, 16GB/节点
Spark3.5.83节点, 8GB/节点5节点, 16GB/节点
Kafka3.7.03 Broker5 Broker
HBase2.5.93节点5节点
MySQL8.04GB8GB

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 2

3. 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. 环境搭建
#

  1. 安装和配置JDK 11、Scala 2.13.8
  2. 部署Hadoop 3.3.6集群
  3. 部署Spark 3.5.8集群
  4. 部署Kafka 3.7.0集群
  5. 部署HBase 2.5.9集群
  6. 安装MySQL 8.0并执行初始化SQL
  7. 配置环境变量(数据库密码等敏感信息)

2. 项目构建
#

  1. 使用Maven构建项目:mvn clean package
  2. 上传jar包到Spark集群
  3. 验证application.conf配置正确性

3. 数据采集
#

  1. 启动Kafka服务和创建主题
  2. 运行SalesDataProducer生成模拟销售数据
  3. 运行SensorDataProducer生成模拟传感器数据
  4. 使用Kafka Console Consumer验证数据

4. 数据存储
#

  1. 运行SensorDataToHBase将传感器数据写入HBase
  2. 运行SalesDataToMySQL将销售数据写入MySQL
  3. 验证数据正确存储
  4. 测试降级策略(模拟HBase/MySQL不可用)

5. 数据分析
#

  1. 运行SalesAnalysis进行销量分析
  2. 运行PriceAnalysis进行价格分析
  3. 运行CustomerPortraitAnalysis进行用户画像分析
  4. 运行SensorAnalysis进行传感器分析
  5. 运行CompetitivenessAnalysis进行竞争力分析

6. 实时处理
#

  1. 运行RealtimeSensorAlert启动实时告警
  2. 验证窗口计算结果
  3. 测试背压机制
  4. 测试降级策略

7. 可视化展示
#

  1. 配置Grafana数据源
  2. 创建仪表板和图表
  3. 验证可视化效果

8. 系统监控
#

  1. 配置Prometheus和Grafana监控
  2. 设置告警规则
  3. 验证监控效果

项目扩展
#

  1. 实时推荐系统:基于客户购买历史和行为数据,使用Spark MLlib实现车辆推荐
  2. 预测分析:使用时间序列模型预测销售趋势,使用分类模型预测传感器故障
  3. 智能诊断:基于OBD诊断数据和传感器数据,实现车辆故障智能诊断
  4. 供应链优化:分析销售数据,使用优化算法优化车辆供应链
  5. 驾驶行为评分:基于GPS轨迹和加速度数据,构建驾驶行为评分模型
  6. 动态定价:基于市场供需和竞争分析,实现动态定价策略

项目总结
#

本项目实现了一个完整的企业级汽车大数据分析系统,涵盖了数据采集、存储、处理、分析和可视化的全流程。通过整合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释放
数据验证空值检查、范围检查、类型检查
异常处理简单printSLF4J日志 + 降级策略 + 重试机制
集合转换JavaConverters(已废弃)scala.jdk.CollectionConverters
降级策略HDFS暂存 + 批处理回退 + 自动恢复
分析维度仅销量和传感器销量/价格/用户画像/竞争力四维指标体系
技术栈版本混乱Spark 3.5.8 + Scala 2.13.8 + JDK 11统一

通过本项目,我们掌握了企业级大数据系统的设计和实现方法,理解了车联网数据特征、Lambda架构、Spark Streaming微批次处理原理以及汽车数据分析指标体系,为未来的大数据项目开发打下了坚实的基础。