一、整体实现流程

这个系统本质上是一个电商用户行为数据分析平台。整体流程是:

用户行为数据 / 商品数据 / 城市数据 → Hive/HDFS存储 → Spark离线计算或Spark Streaming实时计算 → MySQL保存结果 → Vue + Tableau前端展示 → WebSocket异步通知任务状态。

系统分为三层:

数据源层:存储用户行为、用户信息、商品信息、城市信息、广告点击日志。
数据处理层:Spark Core、Spark SQL、Spark Streaming负责离线和实时分析。
数据展示层:MySQL保存分析结果,Vue页面和Tableau图表读取展示。


二、后端实现流程

1. 数据准备与存储

首先准备电商平台的基础数据,包括:

用户信息表 user_info:用户ID、年龄、性别、职业、城市等。
用户行为表 user_action:Session ID、页面ID、行为时间、搜索关键词、点击商品、下单商品、支付商品等。
商品信息表 product_info:商品ID、商品名称、商品类型。
城市信息表 city_info:城市ID、城市名称、所属区域。
任务表 task:保存每次分析任务的参数,比如时间范围、年龄范围、页面流等。

用户行为明细数据存储在 Hive 中,底层依赖 HDFS;任务参数和计算结果存储在 MySQL 中。


2. 任务提交流程

前端页面选择分析功能后,后端会生成一条任务记录,写入 MySQL 的 task 表。

例如页面跳转率任务参数:

{
  "targetPageFlow": "1,2,3,4,5,6,7,8,9",
  "startDate": "2022-04-13",
  "endDate": "2022-04-30"
}

然后后端通过脚本或接口触发对应的 Spark 作业:

spark-submit --class xxx.SessionStatJob ecommerce-analysis.jar 1

其中 1task_id,Spark 作业启动后会根据 task_id 查询 MySQL 中的任务参数。


三、四个核心模块实现

1. 用户访问统计模块

这个模块主要分析用户 Session 的访问质量。(Session :一个用户从进入网站到离开网站期间的一次完整访问过程。)

后端处理流程

第一步,Spark 从 MySQL 的 task 表读取任务参数,比如年龄、性别、城市、时间范围。

第二步,从 Hive 的 user_action 表中读取指定时间范围内的用户行为数据。

第三步,以 session_id 为 key,把用户行为聚合成 Session 粒度数据。

原始数据是一条行为一条记录,例如:

用户A  session_001  点击页面1
用户A  session_001  点击页面2
用户A  session_001  下单商品3

聚合后变成:

session_001 -> 访问时长、访问页面数、搜索关键词、点击品类、下单品类、支付品类

第四步,将 Session 数据和用户信息表 Join,得到带用户属性的 Session 数据。

第五步,根据任务参数过滤数据,比如只保留年龄 10 到 50 岁之间的用户。

第六步,使用 Spark Accumulator 统计访问时长区间和访问步长区间。

例如:

1-3秒
4-6秒
7-9秒
10-30秒
30-60秒
1-3分钟
3-10分钟
10-30分钟
30分钟以上

页面跳转数也按区间统计:

1-3页
4-6页
7-9页
10-30页
30-60页
60页以上

第七步,计算每个区间占总 Session 数的比例,写入 MySQL 的 session_stat 表。

第八步,按照每天每小时的 Session 比例随机抽取 Session,写入 session_random 表。

第九步,从抽样 Session 中统计 Top10 热门品类,排序依据是:

点击数优先,其次下单数,再次支付数

结果写入 top10_category 表。

第十步,统计每个热门品类下点击数最高的 Top10 Session,写入 top10_session 表,同时保存 Session 明细到 session_detail 表。

面试表达

可以这样说:

我负责的用户访问统计模块,核心是把 Hive 中的行为明细数据按照 session_id 聚合成 Session 粒度,然后 Join 用户画像信息,根据任务参数进行过滤。之后用 Spark Accumulator 统计访问时长和访问步长的区间占比,并通过自定义二次排序统计 Top10 热门品类和每个品类下的 Top10 活跃 Session,最终将结果写入 MySQL,供前端展示。


2. 页面跳转率模块

这个模块用来分析用户在页面之间的转化情况。

后端处理流程

第一步,前端传入目标页面流,例如:

1,2,3,4,5,6,7

第二步,后端将页面流拆成页面组合:

1_2
2_3
3_4
4_5
5_6
6_7

第三步,Spark 从 Hive 查询指定时间范围内的用户行为数据。

第四步,按照 session_id 分组。

第五步,对每个 Session 内的行为按照 action_time 排序,得到用户真实访问路径。

例如:

1 -> 2 -> 3 -> 5 -> 7

第六步,将用户访问路径两两切分:

1_2
2_3
3_5
5_7

第七步,判断这些页面组合是否属于目标页面流,如果属于,则累计对应页面组合的访问次数。

第八步,计算跳转率:

1_2跳转率 = 1_2访问次数 / 页面1访问次数
2_3跳转率 = 2_3访问次数 / 1_2访问次数
3_4跳转率 = 3_4访问次数 / 2_3访问次数

第九步,将最终结果写入 MySQL 的 page_convert 表。

面试表达

页面跳转率模块主要是分析用户在指定页面流中的转化情况。我先把前端传入的页面流拆成多个页面 pair,然后将 Hive 中的用户行为按 session_id 分组,并按行为时间排序。接着把每个 Session 的访问路径切分成页面 pair,统计目标页面 pair 的访问次数,最后按漏斗逻辑计算每一步跳转率并写入 MySQL。


3. 热门商品统计模块

这个模块用于统计不同区域的 Top3 热门商品。

后端处理流程

第一步,Spark 从 MySQL 读取任务参数,比如开始日期和结束日期。

第二步,从 Hive 查询这段时间内的用户点击行为。

第三步,从 MySQL 查询城市信息表 city_info

第四步,将用户行为数据和城市信息 Join,得到:

商品ID、城市、区域

第五步,将 Join 后的数据注册成 Spark SQL 临时表。

第六步,按区域和商品分组,统计点击次数。

第七步,使用自定义 UDAF 将城市信息拼接起来。

例如:

华北 商品A 北京:30,天津:20

第八步,Join 商品信息表 product_info,补充商品名称和商品类型。

第九步,使用 Spark SQL 开窗函数:

row_number() over(partition by area order by click_count desc)

对每个区域内的商品点击数进行排序。

第十步,筛选每个区域排名前三的商品,写入 MySQL 的 top3_product 表。

面试表达

热门商品模块主要用 Spark SQL 实现。我先把用户点击行为和城市信息 Join,得到商品在不同区域的点击数据,然后按区域和商品聚合点击次数。为了展示城市分布,我实现了一个 UDAF 函数拼接城市信息。最后使用开窗函数按区域分组排序,取每个区域 Top3 商品,并写入 MySQL。


4. 广告点击实时统计模块

这个模块是实时计算部分,使用 Kafka + Spark Streaming。

后端处理流程

第一步,Kafka Producer 模拟生成广告点击日志。

日志格式可以设计为:

timestamp province city user_id ad_id

例如:

2022-04-20 10:10:01 Beijing Beijing 1001 5

第二步,Spark Streaming 使用 Kafka Direct API 消费 Kafka Topic。

第三步,每隔几秒生成一个 Batch RDD。

第四步,统计用户每天对某个广告的点击次数,更新到 MySQL 的 user_click_count 表。

第五步,如果某个用户一天内对单个广告点击超过 100 次,就将该用户写入 ad_blacklist 黑名单表。

第六步,每个 Batch 处理前,动态读取 MySQL 中的黑名单。

第七步,将当前 Batch 的广告点击数据和黑名单 Join,过滤掉黑名单用户的点击行为。

第八步,使用 updateStateByKey 维护广告点击累计状态。

第九步,统计每天每个省份、城市、广告的点击数,写入 ad_stat 表。

第十步,使用 Spark SQL 开窗函数统计每个省份点击量 Top3 的广告,写入 ad_top3 表。

第十一步,使用窗口函数统计最近一小时内每个广告每分钟的点击趋势,写入 click_trend 表。

面试表达

广告点击模块是实时计算模块。我使用 Kafka 模拟广告点击流,Spark Streaming 通过 Direct API 消费 Kafka 数据。每个 Batch 中先统计用户对广告的点击次数,如果超过阈值就异步更新黑名单。后续 Batch 会动态加载黑名单并过滤异常点击,然后使用 updateStateByKey 维护广告点击状态,最终统计每日广告点击量、省份 Top3 热门广告以及最近一小时点击趋势,结果写入 MySQL 给前端展示。


四、前端实现流程

你的简历写的是 Vue + Tableau + WebSocket,可以这样设计和讲。

1. 前端页面结构

前端主要分为五个页面:

1)首页 Dashboard

展示系统总览:

  • 今日任务数

  • Session 总数

  • 热门品类 Top10

  • 热门商品 Top3

  • 广告点击趋势

  • 黑名单用户数量

2)用户访问统计页面

功能:

  • 选择时间范围

  • 选择年龄范围

  • 选择性别、职业、城市

  • 点击“开始分析”

  • 展示访问时长分布

  • 展示访问步长分布

  • 展示 Top10 热门品类

  • 展示 Top10 活跃 Session

3)页面跳转率页面

功能:

  • 输入目标页面流,例如 1,2,3,4,5,6

  • 选择日期范围

  • 提交分析任务

  • 展示每个页面组合的跳转率

4)区域热门商品页面

功能:

  • 选择日期范围

  • 查看各区域 Top3 商品

  • 使用 Tableau 展示柱状图或地图图表

5)广告实时统计页面

功能:

  • 实时展示广告点击量

  • 展示各省 Top3 广告

  • 展示黑名单用户

  • 展示最近一小时广告点击趋势折线图


2. 前端任务提交流程

用户在 Vue 页面填写筛选条件后,点击“开始分析”。

前端向后端发送请求:

POST /api/task/session

请求体示例:

{
  "startDate": "2022-04-15",
  "endDate": "2022-04-30",
  "startAge": 10,
  "endAge": 50,
  "sex": "male",
  "city": "北京"
}

后端接收到请求后:

  1. 将参数转成 JSON;

  2. 写入 MySQL 的 task 表;

  3. 返回 task_id

  4. 异步提交 Spark 作业;

  5. 前端通过 WebSocket 监听任务状态。


3. WebSocket 异步通信流程

因为 Spark 任务不是瞬间完成的,所以前端不能一直阻塞等待。

流程是:

  1. 前端提交任务;

  2. 后端立即返回 task_id

  3. 前端建立 WebSocket 连接;

  4. 后端任务状态变化时主动推送;

  5. 前端收到完成通知后,请求结果接口;

  6. 前端刷新图表。

例如后端推送:

{
  "taskId": 1,
  "status": "FINISHED",
  "message": "用户访问统计任务执行完成"
}

前端收到后调用:

GET /api/result/session/1

然后将数据渲染到页面。


4. Tableau 可视化流程

Vue 页面中可以通过 iframe 或 Tableau JS API 嵌入 Tableau 图表。

数据来源是 MySQL 结果表,比如:

  • session_stat

  • top10_category

  • page_convert

  • top3_product

  • ad_top3

  • click_trend

Tableau 连接 MySQL 后,根据结果表生成柱状图、折线图、漏斗图等。

Vue 页面负责:

  • 任务提交;

  • 参数选择;

  • 页面路由;

  • 嵌入 Tableau 图表;

  • 通过 WebSocket 刷新图表。


五、可以写进简历/面试的完整项目流程

你可以这样讲:

这个项目是一个基于 Spark 的电商用户行为分析系统。我主要负责前后端开发和数据分析任务实现。系统底层使用 Hadoop/HDFS 和 Hive 存储用户行为数据,Kafka 模拟实时广告点击流,Spark Core 和 Spark SQL 负责离线分析,Spark Streaming 负责实时计算,MySQL 保存任务参数和分析结果,前端使用 Vue 和 Tableau 展示统计图表,并通过 WebSocket 实现任务状态异步通知。

后端流程是:前端提交分析条件后,后端将任务参数写入 MySQL 的 task 表,然后通过 spark-submit 提交对应 Spark 作业。Spark 作业根据 task_id 读取任务参数,从 Hive 或 Kafka 中读取数据,完成 Session 聚合、页面跳转率计算、区域热门商品排序和广告点击实时统计,最后将结果写入 MySQL。前端再通过接口查询结果,并使用 Tableau 图表进行可视化展示。

我实现的核心模块包括:用户访问统计、页面跳转率分析、区域 Top3 热门商品统计和广告点击实时分析。其中用户访问统计使用 Spark Core 对 Session 聚合,并通过 Accumulator 统计访问时长和访问步长区间;页面跳转率模块按 Session 内行为时间排序后计算页面流转化率;热门商品模块使用 Spark SQL、UDAF 和开窗函数统计各区域 Top3 商品;广告点击模块使用 Spark Streaming 结合 Kafka Direct API 消费实时点击流,通过 StateByKey 维护点击状态,并异步更新黑名单。


Spark Standalone 是 Spark 自己管资源;Kubernetes 是 K8s 管容器,Spark 跑在容器里。

1. Spark Standalone 怎么做

你的原项目更接近这个。

架构

Spark Master
    ↓
Spark Worker 1
Spark Worker 2
Spark Worker 3

运行流程

spark-submit
→ 提交 Jar 到 Spark Master
→ Master 分配 Worker 资源
→ Worker 启动 Executor
→ Executor 执行 Spark Task
→ 结果写入 MySQL

实现步骤

  1. 三台虚拟机安装 JDK、Scala、Spark;

  2. 配置 spark-env.sh

  3. 配置 slaves/workers

  4. 启动集群:

start-master.sh
start-workers.sh
  1. 提交任务:

spark-submit \
  --class com.xxx.SessionStatApp \
  --master spark://spark1:7077 \
  ecommerce-analysis.jar \
  1

这里的 1 就是 task_id


2. Kubernetes 怎么做

K8s 版本是把 Spark 任务容器化。

架构

Kubernetes Master
    ↓
Node 1: Driver Pod
Node 2: Executor Pod
Node 3: Executor Pod

运行流程

spark-submit
→ 请求 Kubernetes API
→ 创建 Driver Pod
→ Driver Pod 创建 Executor Pod
→ Executor Pod 执行 Spark Task
→ 结果写入 MySQL

实现步骤

第一步,把 Spark 作业做成 Docker 镜像:

FROM apache/spark:3.5.0

COPY ecommerce-analysis.jar /opt/app/ecommerce-analysis.jar
COPY core-site.xml /opt/spark/conf/
COPY hdfs-site.xml /opt/spark/conf/
COPY hive-site.xml /opt/spark/conf/

第二步,把镜像推到镜像仓库:

docker build -t registry.xxx.com/ecommerce-analysis:1.0 .
docker push registry.xxx.com/ecommerce-analysis:1.0

第三步,提交 Spark on K8s 任务:

spark-submit \
  --master k8s://https://k8s-apiserver:6443 \
  --deploy-mode cluster \
  --name ecommerce-session-job \
  --class com.xxx.SessionStatApp \
  --conf spark.kubernetes.container.image=registry.xxx.com/ecommerce-analysis:1.0 \
  --conf spark.executor.instances=3 \
  --conf spark.executor.memory=2g \
  --conf spark.executor.cores=1 \
  local:///opt/app/ecommerce-analysis.jar \
  1

3. 核心区别

对比点 Spark Standalone Kubernetes
调度对象 Driver、Executor Pod、容器
部署方式 直接装在虚拟机上 Docker 镜像部署
资源管理 Spark Master 管 K8s Scheduler 管
扩缩容 手动加 Worker 可以弹性扩缩容
运维复杂度
适合场景 学习、毕设、小集群 生产环境、云原生
和你项目匹配度 很匹配 能用但偏重

4. 面试怎么讲最稳

你可以这样说:

原项目采用的是 Spark Standalone 模式,Spark Master 负责接收 spark-submit 提交的任务,并将 Executor 分配到 Worker 节点上执行。
如果迁移到 Kubernetes,整体计算逻辑不用变,主要变化在部署层:需要将 Spark 作业、JDK、依赖包和配置文件打成 Docker 镜像,然后通过 spark-submit 提交到 K8s API Server,由 Kubernetes 创建 Driver Pod 和 Executor Pod。K8s 负责 Pod 调度和资源隔离,Spark Driver 仍然负责任务 DAG 切分和 Task 调度。

这句很关键:

K8s 只替代资源部署和容器调度,不替代 Spark 的计算逻辑。

怎么使用Kubernetes管理集群

我们主要是把 Spark 分析任务容器化后提交到 Kubernetes 集群运行。Kubernetes 不直接参与 Spark 的计算逻辑,它主要负责 Driver Pod 和 Executor Pod 的创建、资源分配、调度、失败重启和生命周期管理。Spark 内部仍然负责任务 DAG 切分、Stage 划分和 Task 调度。

具体实现可以这样讲:

前端提交分析任务
→ 后端写入 task 表
→ 后端调用 spark-submit
→ spark-submit 提交到 K8s API Server
→ K8s 创建 Spark Driver Pod
→ Driver 读取 MySQL task 参数
→ Driver 申请 Executor Pod
→ Executor 从 Hive/HDFS 或 Kafka 读取数据
→ Spark 执行计算
→ 结果写回 MySQL
→ 前端查询结果并展示

然后展开:

首先我把 Spark 作业、JDK、依赖包、项目 Jar、Hadoop/Hive 配置文件打成 Docker 镜像。镜像里包含 core-site.xmlhdfs-site.xmlhive-site.xml,保证 Pod 启动后能访问 HDFS 和 Hive。

然后通过 spark-submit 指定 Kubernetes 作为 master,例如:

spark-submit \
  --master k8s://https://k8s-apiserver:6443 \
  --deploy-mode cluster \
  --name ecommerce-session-job \
  --class com.xxx.SessionStatApp \
  --conf spark.kubernetes.container.image=registry.xxx.com/ecommerce-analysis:1.0 \
  --conf spark.executor.instances=3 \
  --conf spark.executor.memory=2g \
  --conf spark.executor.cores=1 \
  local:///opt/app/ecommerce-analysis.jar \
  1

这里 1 是任务 ID。Driver Pod 启动后,会根据 task_id 从 MySQL 读取任务参数,然后通过 SparkContext 构建 DAG。Kubernetes 根据资源配置创建多个 Executor Pod,每个 Executor Pod 负责执行 Spark Task。任务执行完成后,结果写入 MySQL,Driver Pod 结束,Executor Pod 也释放。

如果问“调度具体怎么调”,你说:

Kubernetes 调度的是 Pod,不是 Spark 的 Task。K8s Scheduler 会根据每个节点的 CPU、内存资源,把 Driver Pod 和 Executor Pod 调度到合适的 Node 上。Spark 的 Driver 再把具体的 Task 分发给 Executor 执行。所以这里是两层调度:K8s 负责容器级调度,Spark 负责计算任务级调度。

如果问“你用了哪些 K8s 对象”,你说:

主要涉及 Pod、ConfigMap、Secret、ServiceAccount。
ConfigMap 用来挂载 Hadoop/Hive 配置文件;Secret 用来保存 MySQL 账号密码;ServiceAccount 用来给 Driver Pod 创建 Executor Pod 的权限;Pod 是 Driver 和 Executor 的运行载体。Kafka、MySQL、HDFS 可以部署在 K8s 外部,通过 Service 地址或外部域名访问。

最稳总结:

所以 Kubernetes 在这个项目里的作用不是替代 Spark,而是把 Spark 作业容器化运行,统一管理 Driver 和 Executor 的资源、调度和生命周期。计算逻辑仍然由 Spark 完成。


Kubernetes和Spark调度的区别

一、Pod 和 Task 根本不是一个层级

它们分别属于:

层级 谁负责 干什么
Kubernetes Pod调度 调度“运行环境”
Spark Task调度 调度“计算任务”

所以:

K8s 管“容器在哪台机器运行”
Spark 管“数据怎么计算”

二、你可以把它理解成“两层调度”

第一层:Kubernetes 调度 Pod

K8s Scheduler 决定:

Driver Pod 放哪台机器
Executor Pod 放哪台机器

例如:

Node1:
  Driver Pod

Node2:
  Executor Pod 1

Node3:
  Executor Pod 2

K8s 只关心:

CPU够不够
内存够不够
节点是否健康

它根本不懂:

Session聚合
Top10商品
RDD

这些它完全不知道。


第二层:Spark 调度 Task

真正的数据计算:

RDD切分
Stage划分
Task分发
Shuffle

这些都是:

Spark Driver

负责。


三、真正执行流程(非常重要)

你这个问题其实是 Spark on K8s 的核心。


第一步:提交任务

spark-submit ...

第二步:K8s 创建 Driver Pod

例如:

Node1:
  spark-driver-pod

第三步:Driver Pod 启动 Spark Driver

Driver 开始:

读取Hive
构建DAG
划分Stage
划分Task

第四步:Driver 向 K8s 申请 Executor Pod

例如:

申请3个Executor

K8s 决定:

Executor1 -> Node2
Executor2 -> Node3
Executor3 -> Node2

第五步:Executor Pod 启动

每个 Pod 里运行:

Spark Executor JVM

第六步:Spark Driver 分发 Task

例如:

Task1 -> Executor1
Task2 -> Executor2
Task3 -> Executor3

这里:

不是K8s分Task

而是:

Spark Driver分Task

这是最关键的点。


四、Task 到底是什么

Task 本质:

对一个数据分区的计算

例如:

你有:

100GB 用户行为数据

Spark 切成:

200个Partition

那么:

每个Partition对应一个Task

比如:

Task1 处理 partition1
Task2 处理 partition2
Task3 处理 partition3

五、Executor 和 Task 的关系

Executor:

真正执行Task的进程

一个 Executor 可以同时跑多个 Task。

例如:

Executor:
  4 cores

那么:

最多同时跑4个Task

六、你刚才那个理解差在哪

你刚才说:

K8s负责调度不同服务器,然后spark负责调度同一服务器上的任务?

这个:

前半句对

K8s负责Pod在哪台机器运行

正确。


后半句不完全对

Spark:

不是只调度“同一服务器”

而是:

调度整个集群上的Task

例如:

Node2 Executor
Node3 Executor

Spark Driver 都会统一调度。


七、真正关系图(面试非常加分)

                Kubernetes
        (Pod/容器级调度)
                       ↓
     Driver Pod                Executor Pod
                                  Executor Pod
                                  Executor Pod
                       ↓
                  Spark Driver
             (Task级调度)
                       ↓
            Task1 Task2 Task3 ...

八、最标准面试回答(建议背)

你可以直接这样说:

Kubernetes 和 Spark 属于两层调度。Kubernetes 负责 Pod 的资源调度和节点分配,比如决定 Driver Pod 和 Executor Pod 运行在哪台 Node 上;而 Spark Driver 负责计算层面的 DAG 划分、Stage 划分和 Task 调度。Task 本质是对一个数据 Partition 的计算任务,Spark Driver 会把 Task 分发给不同 Executor 执行,所以 K8s 调度的是运行环境,Spark 调度的是真正的数据计算任务。


driver和executor pod分别干什么的,怎么用
spark里的SparkContext 、DAG、RDD、Stage、shuffle分别是什么和什么关系

一、Driver Pod 和 Executor Pod 是什么

这是 Spark on Kubernetes 最核心的东西。

可以理解成:

组件 类比 干什么
Driver Pod 项目经理 负责“指挥”
Executor Pod 工人 负责“干活”

1. Driver Pod 干什么

Driver Pod 里运行的是:

Spark Driver 进程

它是整个 Spark 作业的大脑。

负责:

1. 创建 SparkContext
2. 读取任务参数
3. 构建 DAG
4. 划分 Stage
5. 申请 Executor
6. 分发 Task
7. 汇总结果

例如你项目里的:

统计 Top10 商品

Driver 会:

读取 Hive 数据
→ 分析 map/reduce/groupBy
→ 切分执行阶段
→ 把任务发给 Executor

所以:

Driver 不负责真正算数据

它主要:

负责调度

2. Executor Pod 干什么

Executor Pod 是真正干活的。

每个 Executor Pod 里运行:

Spark Executor JVM

负责:

1. 执行 Task
2. 读写数据
3. shuffle 数据交换
4. 缓存RDD
5. 返回结果给Driver

例如:

100GB用户行为数据

被切成:

200个Partition

那么:

Executor1:
  Task1
  Task2

Executor2:
  Task3
  Task4

真正计算:

map
filter
reduceByKey
groupByKey

都在 Executor 里执行。


3. Spark on K8s 怎么用


提交任务

spark-submit \
  --master k8s://https://k8s-apiserver:6443 \
  --deploy-mode cluster \
  --class com.xxx.SessionApp \
  --conf spark.executor.instances=3 \
  local:///opt/app.jar

执行流程

第一步

K8s 创建:

Driver Pod

第二步

Driver Pod 启动:

SparkContext

第三步

Driver 向 K8s 申请:

Executor Pod

第四步

K8s 调度:

Executor Pod

到不同 Node。


第五步

Driver:

分发Task

Executor:

执行Task

二、SparkContext 是什么

SparkContext:

Spark程序入口

类似:

SpringBoot里的SpringApplication

创建:

SparkConf conf = new SparkConf()
        .setAppName("session")
        .setMaster("local");

SparkContext sc = new SparkContext(conf);

它负责:

连接集群
创建RDD
调度任务

SparkContext 可以理解成:

Spark Driver 和集群通信的总控制器

三、RDD 是什么

RDD:

Resilient Distributed Dataset
弹性分布式数据集

本质:

Spark里的“分布式数据”

例如:

JavaRDD<String> lines =
    sc.textFile("hdfs://user_action.log");

这里:

lines

就是:

RDD

RDD 特点:

特点 含义
分布式 数据分散在多机器
不可变 只能生成新RDD
可并行 多机器同时算
容错 丢了能重算

四、DAG 是什么

DAG:

Directed Acyclic Graph
有向无环图

本质:

Spark的计算流程图

例如:

rdd.map()
   .filter()
   .groupByKey()
   .reduceByKey()

Spark 不会立刻执行。

而是先生成:

DAG执行图

类似:

读取数据
   ↓
map
   ↓
filter
   ↓
groupByKey
   ↓
reduceByKey

DAG 是:

Spark Driver生成的执行计划

五、Stage 是什么

Stage:

Spark执行阶段

Spark 会把 DAG:

切成多个 Stage

切分依据:

是否发生 Shuffle

例如:

rdd.map()
   .filter()
   .groupByKey()

会变成:


Stage1

map
filter

Shuffle

数据重新分区


Stage2

groupByKey

六、Shuffle 是什么(最核心)

Shuffle:

跨节点数据重分布

这是 Spark 最重的操作。


例如:

原数据:

Node1:
  a 1
  b 1

Node2:
  a 1
  c 1

现在:

reduceByKey(_ + _)

必须:

把相同key放一起

所以:

a的数据
→ 同一个Executor

这就是:

Shuffle

Shuffle 会:

网络传输
磁盘IO
序列化
反序列化

所以:

很慢

七、它们之间的关系(超级重要)

真正执行流程:

SparkContext
    ↓
创建RDD
    ↓
RDD转换生成DAG
    ↓
DAG按Shuffle切分Stage
    ↓
Stage切成多个Task
    ↓
Task分发给Executor
    ↓
Executor执行计算

八、举一个你项目里的真实例子

例如:

rdd.mapToPair()
   .groupByKey()
   .reduceByKey()

Spark内部流程


1. SparkContext 创建

Driver启动。


2. 读取Hive数据生成RDD

user_action_rdd

3. 构建DAG

mapToPair
↓
groupByKey
↓
reduceByKey

4. 遇到 Shuffle

groupByKey

触发 Shuffle。


5. DAG 划分 Stage

Stage1:
  mapToPair

Shuffle

Stage2:
  groupByKey
  reduceByKey

6. Stage 划分 Task

假设:

200 partitions

则:

200 Task

7. Driver 分发给 Executor

Executor1:
  Task1~50

Executor2:
  Task51~100

九、面试里最标准回答(建议背)

SparkContext 是 Spark 程序入口,负责与集群通信和创建 RDD。RDD 是 Spark 的分布式数据抽象,所有数据计算都会基于 RDD 转换生成 DAG 执行图。Spark Driver 会根据 DAG 中是否发生 Shuffle 将计算切分成多个 Stage,每个 Stage 再按 Partition 划分成多个 Task。Task 最终会分发给不同 Executor 执行。Shuffle 本质是跨节点的数据重新分布,比如 groupByKey、reduceByKey 都会触发 Shuffle,也是 Spark 中开销最大的部分之一。


Hadoop和Hive的关系:大数据底层平台和数据仓库管理工具

Hadoop 是“大数据底层平台”
Hive 是“构建在 Hadoop 上的数据仓库工具”

或者更直白:

Hadoop 负责存
Hive 负责查

你项目里的真实关系

你项目实际上是:

用户行为日志
→ 存到 HDFS(Hadoop)
→ Hive 建表
→ Spark SQL 查询 Hive

一、Hadoop 到底是什么

Hadoop 不是一个单独软件。

它是一套生态。

最核心的是:

组件 作用
HDFS 分布式存储
MapReduce 分布式计算
Yarn 资源调度

你项目里真正用到最多的是:

HDFS

它负责:

存储海量日志文件

例如:

user_action.log
ad_click.log

这些文件会切成 block 分布在多个机器上。


二、Hive 是什么

Hive 本质是:

把 HDFS 上的文件“变成数据库表”

因为:

直接处理 HDFS 文件太痛苦

比如原始文件:

2022-04-01,user1,session1,page1,click
2022-04-01,user1,session1,page2,search

如果不用 Hive:

你得:

自己写代码解析文件
自己切字段
自己统计

很麻烦。


所以 Hive 干了件事:

给 HDFS 文件套一层 SQL

三、Hive 和 MySQL 最大区别

很多人容易混。


MySQL

是真数据库:

数据真正存数据库里

Hive

不是传统数据库。

Hive:

只是“元数据管理 + SQL解析”

真正的数据:

仍然存在 HDFS

Hive 只记录:

这个表对应哪个HDFS目录
字段是什么
分隔符是什么

举个你项目里的真实例子


第一步:日志存在 HDFS

例如:

hdfs://spark1:9000/user/hive/warehouse/user_action/

里面有:

part-00001
part-00002

第二步:Hive 建表

CREATE EXTERNAL TABLE user_action (
  date STRING,
  user_id BIGINT,
  session_id STRING,
  page_id BIGINT,
  action_type STRING
)
ROW FORMAT DELIMITED
FIELDS TERMINATED BY ',';

这时候:

Hive 只是记录:

这个表对应哪个HDFS目录

数据没复制。


第三步:Spark SQL 查询 Hive

sqlContext.sql("select * from user_action")

实际上:

Spark
→ Hive元数据
→ 找到HDFS路径
→ 读取HDFS文件

所以关系图其实是

        Hive
          ↓
SQL解析 + 元数据
          ↓
        HDFS
(真正存数据)

四、为什么大数据项目都喜欢 Hive

因为:

没 Hive 时

你需要:

Java代码
自己读文件
自己split
自己解析

有 Hive 后

直接:

select *
from user_action
where date='2022-04-01'

简单太多。


五、你项目里 Hadoop 和 Hive 的分工

Hadoop(HDFS)

负责:

海量用户行为日志存储

Hive

负责:

把日志映射成结构化表

Spark

负责:

真正做计算分析

六、面试里最标准回答

你可以这样说:

Hadoop 在项目中主要负责底层分布式存储,用户行为日志会存储在 HDFS 中;Hive 构建在 Hadoop 之上,负责将 HDFS 上的日志文件映射成结构化表,方便使用 SQL 查询;Spark 再通过 Spark SQL 读取 Hive 表进行离线分析和实时统计。Hive 本身并不真正存储数据,它更像是 Hadoop 上的数据仓库和 SQL 查询层。

Logo

电商企业物流数字化转型必备!快递鸟 API 接口,72 小时快速完成物流系统集成。全流程实战1V1指导,营造开放的API技术生态圈。

更多推荐