电商用户行为分析系统
一、整体实现流程
这个系统本质上是一个电商用户行为数据分析平台。整体流程是:
用户行为数据 / 商品数据 / 城市数据 → 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
其中 1 是 task_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": "北京"
}
后端接收到请求后:
-
将参数转成 JSON;
-
写入 MySQL 的
task表; -
返回
task_id; -
异步提交 Spark 作业;
-
前端通过 WebSocket 监听任务状态。
3. WebSocket 异步通信流程
因为 Spark 任务不是瞬间完成的,所以前端不能一直阻塞等待。
流程是:
-
前端提交任务;
-
后端立即返回
task_id; -
前端建立 WebSocket 连接;
-
后端任务状态变化时主动推送;
-
前端收到完成通知后,请求结果接口;
-
前端刷新图表。
例如后端推送:
{
"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
实现步骤
-
三台虚拟机安装 JDK、Scala、Spark;
-
配置
spark-env.sh; -
配置
slaves/workers; -
启动集群:
start-master.sh
start-workers.sh
-
提交任务:
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.xml、hdfs-site.xml、hive-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 查询层。
更多推荐



所有评论(0)