一、Map Task 与 Reduce Task的作用

1.1 Map Task和Reduce Task

1.1.1 Map Task(映射任务)
  • 定义:处理输入数据分区的初始任务,执行mapfilterflatMap窄依赖转换操作。 【详细的宽窄依赖列举见第三点】

  • 特点

    • 每个Map Task独立处理一个输入分区
    • 输出结果直接传递给下游任务或写入本地磁盘(Shuffle Write)
    • 在Spark中对应ShuffleMapTask类型
  • 示例

    // Map Task处理逻辑
    val mapped = inputRDD.map(x => (x.key, x.value))  // 每个分区独立处理
    
1.1.2 Reduce Task(归约任务)
  • 定义:从多个Map Task收集数据并进行聚合的任务,执行reduceByKeygroupByKeyjoin宽依赖操作。

  • 特点

    • 需要从多个Map Task拉取数据(Shuffle Read)
    • 对相同key的数据进行聚合计算
    • 在Spark中通常对应ResultTask类型
  • 示例

    // Reduce Task处理逻辑
    val reduced = mapped.reduceByKey(_ + _)  // 需要跨分区聚合
    

区别

特性Map TaskReduce Task
数据依赖窄依赖(分区内处理)宽依赖(跨分区聚合)
数据来源输入数据或上游RDD多个Map Task的输出
网络传输通常无(除非数据本地性不足)必须从多个节点拉取数据
输出目标下游任务或本地磁盘最终结果或下一阶段
任务类型ShuffleMapTaskResultTask

Stage 1: Reduce Tasks

网络Shuffle传输

Stage 0: Map Tasks

Executor 5

Executor 4

Executor 3

Executor 2

Executor 1

网络传输

网络传输

网络传输

分区1数据

Map Task

Shuffle Write
分区排序

本地磁盘

分区2数据

Map Task

Shuffle Write
分区排序

本地磁盘

分区3数据

Map Task

Shuffle Write
分区排序

本地磁盘

Shuffle Manager

Shuffle Read
获取对应分区

Reduce Task

结果聚合

Shuffle Read
获取对应分区

Reduce Task

结果聚合

最终结果

注意点:

  1. stage0阶段执行Map Task任务;stage1阶段执行Reduce Task任务;

  2. Map阶段和Reduce阶段中,有多个Executor并行执行。

1.2 Spark的Stage1对比Stage2

特性Stage 0(Map阶段)Stage 1(Reduce阶段)
数据依赖无,读取输入数据或上一个RDD依赖Stage 0的输出,需要shuffle read
网络传输通常无(除非数据本地性不足)需要从多个节点拉取数据
输出写入本地磁盘(shuffle write)输出最终结果或传递给下一个操作
任务类型通常是Map任务通常是Reduce任务(但Spark中不严格区分)
1.2.1 Stage 0(Map阶段)
// 示例:Stage 0的操作
val rdd = sc.textFile("data.txt")          // 读取数据
  .flatMap(_.split(" "))                   // 转换操作1
  .map(word => (word, 1))                  // 转换操作2
  // 这些操作都在Stage 0中,因为它们之间是窄依赖

特点

  • 窄依赖:每个父RDD的分区最多被子RDD的一个分区使用
  • 无需Shuffle:数据在同一个分区内处理
  • 流水线优化:多个转换操作可以合并执行
1.2.2 Stage 1(Reduce阶段)
// 示例:Stage 1的操作
val wordCounts = mappedRDD.reduceByKey(_ + _)  // 触发新的Stage

二、DataFrame操作的依赖类型

2.1表格汇总

操作类型具体操作依赖类型是否触发Action是否触发Shuffle说明
转换操作select()窄依赖列投影,不改变数据分区
filter() / where()窄依赖行过滤,不改变数据分区
withColumn()窄依赖添加/修改列,不改变分区
withColumnRenamed()窄依赖重命名列,不改变分区
drop()窄依赖删除列,不改变分区
limit()窄依赖限制行数,不改变分区
union() / unionAll()窄依赖合并DataFrame,分区数相加
intersect()宽依赖求交集,需要去重
except()宽依赖求差集,需要去重
distinct()宽依赖去重操作,需要全局比较
dropDuplicates()宽依赖删除重复行,需要Shuffle
聚合操作groupBy() + 聚合函数宽依赖分组聚合必须Shuffle
rollup()宽依赖多维聚合,需要Shuffle
cube()宽依赖全维度聚合,需要Shuffle
pivot()宽依赖数据透视,需要Shuffle
排序操作orderBy() / sort()宽依赖全局排序,需要Shuffle
sortWithinPartitions()窄依赖分区内排序,不Shuffle
重分区操作repartition()宽依赖增加/重分布分区
coalesce()窄依赖减少分区,不触发Shuffle
repartitionByRange()宽依赖按范围重分区
Join操作join() (非广播)宽依赖常规Join需要Shuffle
join() (广播)窄依赖广播小表,无Shuffle
crossJoin()宽依赖笛卡尔积,需要Shuffle
窗口函数窗口函数 + partitionBy()宽依赖窗口函数通常需要Shuffle
窗口函数 (无partitionBy)宽依赖全局窗口需要Shuffle
其他操作sample()窄依赖抽样,不改变分区
randomSplit()窄依赖随机分割,不Shuffle
cache() / persist()窄依赖缓存操作,不Shuffle
repartition(col)宽依赖按列重分区
— 以下是新增的Action操作 (会立即触发计算) —
行动操作count()不适用 (Action)可能返回数据集总行数。若上游有Shuffle操作则触发。
collect()不适用 (Action)可能将所有数据拉取到Driver端。若上游有Shuffle操作则触发。
show()不适用 (Action)可能打印前N行数据。若上游有Shuffle操作则触发。
take() / head()不适用 (Action)可能取前N条数据到Driver端。若上游有Shuffle操作则触发。
first()不适用 (Action)可能取第一条数据。等价于 take(1)
foreach()不适用 (Action)可能对每条数据应用函数。若上游有Shuffle操作则触发。
save() / write()不适用 (Action)可能将数据写入外部存储。若上游有Shuffle操作则触发。
saveAsTable()不适用 (Action)可能将数据保存为表。若上游有Shuffle操作则触发。
  • ⏳ (Transformation):惰性操作,不立即计算。
  • ✅ (Action):立即触发计算。
  • 是否触发Shuffle:对于Transformation,标记其本身特性;对于Action,标记其可能因触发包含宽依赖的DAG而导致Shuffle。

2.2 宽窄依赖分析例子

窄依赖操作示例
# 窄依赖操作链 - 无Shuffle,在一个Stage内执行
df_transformed = (df
    .select("id", "name", "salary")
    .filter(col("salary") > 5000)
    .withColumn("bonus", col("salary") * 0.1)
    .withColumnRenamed("id", "employee_id"))
宽依赖操作示例
# 宽依赖操作 - 触发Shuffle,创建新Stage
df_aggregated = (df
    .groupBy("department")
    .agg(
        avg("salary").alias("avg_salary"),
        count("*").alias("employee_count")
    )
    .orderBy("avg_salary"))
混合依赖示例
# 混合依赖 - 会创建多个Stage
result = (df
    # Stage 1: 窄依赖操作
    .select("dept", "name", "salary", "year")
    .filter(col("year") == 2023)
    
    # Stage 2: 宽依赖 (groupBy触发Shuffle)
    .groupBy("dept")
    .agg(avg("salary").alias("avg_salary"))
    
    # Stage 3: 宽依赖 (join触发Shuffle)
    .join(dept_info, "dept", "inner")
    
    # Stage 4: 窄依赖操作
    .select("dept", "dept_name", "avg_salary")
    .orderBy("avg_salary"))

三、Map Task 与 Reduce Task的案例分析

示例:统计海量文本中每个单词的出现次数

1. 输入数据

假设我们有三个文本片段(实际中可能是TB级数据):

文本1: "hello world hello"
文本2: "hello mapreduce"
文本3: "world hello bigdata"
2. MapReduce 执行流程
步骤1:分割(Split)
  • 系统自动将输入文件分割成固定大小的分片(例如128MB),每个分片由一个 Map 任务 处理。
  • 本例中,假设每个文本作为一个分片,共3个分片。
步骤2:Map 阶段(映射)
  • 每个 Map 任务读取一个分片,并逐行处理。

  • Map 函数接收每一行文本,将其拆分为单词,并为每个单词输出 中间键值对<单词, 1>

  • Map 函数伪代码

    def map(line):
        for word in line.split():
            emit(word, 1)  # 输出 (word, 1)
    
  • Map 输出结果

    • Map1(处理文本1):(hello,1), (world,1), (hello,1)
    • Map2(处理文本2):(hello,1), (mapreduce,1)
    • Map3(处理文本3):(world,1), (hello,1), (bigdata,1)
步骤3:Shuffle(洗牌)与排序
  • 系统自动将 所有 Map 输出相同键 的键值对分组,发送到同一个 Reduce 任务。
  • 分组结果:
    • hello → [1, 1, 1, 1](来自 Map1、Map2、Map3)
    • world → [1, 1]
    • mapreduce → [1]
    • bigdata → [1]
步骤4:Reduce 阶段(归约)
  • 每个 Reduce 任务接收一个键及其对应的值列表,进行聚合计算(如求和)。

  • Reduce 函数伪代码

    def reduce(word, counts):
        total = sum(counts)  # 对计数列表求和
        emit(word, total)
    
  • Reduce 输出结果

    • Reduce1 处理 hello(hello, 4)
    • Reduce2 处理 world(world, 2)
    • Reduce3 处理 mapreduce(mapreduce, 1)
    • Reduce4 处理 bigdata(bigdata, 1)
步骤5:输出结果
  • 所有 Reduce 任务的输出合并为最终结果,写入分布式文件系统(如 HDFS):
hello 4
world 2
mapreduce 1
bigdata 1

3. 图示流程

原始数据(分布到3个节点上):
  文本1 → "hello world hello"
  文本2 → "hello mapreduce"
  文本3 → "world hello bigdata"

Map阶段(并行处理):
  Map1 → (hello,1), (world,1), (hello,1)
  Map2 → (hello,1), (mapreduce,1)
  Map3 → (world,1), (hello,1), (bigdata,1)

Shuffle阶段(自动分组):
  hello → [1,1,1,1]
  world → [1,1]
  mapreduce → [1]
  bigdata → [1]

Reduce阶段(并行聚合):
  Reduce1 (hello) → 求和得4 → (hello,4)
  Reduce2 (world) → 求和得2 → (world,2)
  Reduce3 (mapreduce) → 求和得1 → (mapreduce,1)
  Reduce4 (bigdata) → 求和得1 → (bigdata,1)

最终输出:
  hello 4
  world 2
  mapreduce 1
  bigdata 1

更多推荐