本文还有配套的精品资源,点击获取 menu-r.4af5f7ec.gif

简介:Spark 3.1.1 是 Apache Spark 的一个重要版本,结合 Hadoop 2.7 的分布式存储与资源管理能力,构建了高性能、高可靠的大数据处理环境。本文深入解析 Spark 在 Hadoop 平台上的核心特性,涵盖 SQL 性能优化、内存计算、实时流处理、机器学习支持及容错机制等关键内容。通过 Spark 与 YARN、HDFS 的深度集成,实现资源高效调度与大规模数据处理,适用于数据分析、实时计算和 AI 工作负载。本项目可作为企业级大数据平台搭建的参考实践。
Spark

1. Spark 3.1.1 核心架构与内存计算模型解析

Spark核心架构概述

Apache Spark 3.1.1 采用“分层抽象、统一引擎”设计思想,构建于JVM之上,其核心架构由 Driver进程、Executor进程、Cluster Manager与DAGScheduler 等关键组件构成。Driver负责任务调度与上下文管理,Executor执行具体任务并缓存数据,通过 RDD(弹性分布式数据集) 实现不可变、分区、容错的并行计算模型。

内存计算模型与执行优化

Spark将数据优先加载至内存进行迭代计算,显著提升处理效率。其内存管理采用 统一内存池模型(Unified Memory Management) ,动态分配执行(shuffle、join)与存储(缓存)内存。配合 Tungsten二进制存储格式 ,减少对象开销,提升序列化速度与缓存效率。

// 示例:RDD内存缓存操作
val rdd = spark.sparkContext.textFile("hdfs://data.txt")
  .map(_.split(","))
  .cache() // 显式缓存到内存,供多次复用

cache() 方法底层调用 persist(StorageLevel.MEMORY_ONLY) ,将分区数据以反序列化形式驻留内存,加速后续迭代或重复计算场景。

2. Catalyst优化器与SQL执行引擎深度剖析

Apache Spark自2.0版本引入统一的SQL执行引擎以来,其核心组件Catalyst优化器便成为支撑高性能查询处理的关键技术。不同于传统数据库中基于规则或成本的独立优化方式,Catalyst结合了函数式编程范式、树形结构变换和运行时代码生成等现代软件工程理念,构建了一套灵活且高效的查询优化框架。该优化器不仅支持标准SQL语义解析与逻辑计划优化,还能根据底层数据特征动态调整物理执行策略,并通过Whole-Stage Code Generation大幅提升表达式求值效率。在实际生产环境中,理解Catalyst的工作机制对于诊断慢查询、调优执行计划以及提升整体作业性能具有重要意义。

本章将深入解析Catalyst优化器的设计哲学与实现细节,从理论基础出发,逐步展开到查询计划的生成流程、运行时代码生成机制,最后结合真实场景进行性能调优案例分析。通过对这一系列关键技术点的层层拆解,读者将建立起对Spark SQL执行引擎的系统性认知,掌握如何借助执行计划可视化工具和统计信息驱动优化决策的能力。

2.1 Catalyst优化器的理论基础

Catalyst优化器之所以能够在复杂查询场景下保持高灵活性与可扩展性,根本原因在于其建立在坚实的理论基础之上——即以不可变树结构为核心的数据表示模型,配合函数式编程思想下的递归变换机制。这种设计使得整个查询优化过程具备良好的模块化特性,便于开发者添加新的优化规则或适配新型数据源。更重要的是,它实现了规则驱动(Rule-based)与成本驱动(Cost-based)两种优化策略的有机融合,从而兼顾了通用性与精确性。

2.1.1 树结构与函数式编程在查询优化中的应用

在Catalyst中,无论是SQL语句的抽象语法树(AST),还是后续生成的逻辑计划(Logical Plan)和物理计划(Physical Plan),都被建模为 不可变的树形结构 。每个节点代表一个操作符(如Filter、Project、Join等),并通过子节点形成层次化依赖关系。这种设计天然契合函数式编程中的代数数据类型(Algebraic Data Types, ADT)概念,在Scala语言中可通过 case class 和模式匹配(Pattern Matching)高效实现。

例如,一个简单的SQL查询:

SELECT name FROM users WHERE age > 30;

会被解析为如下树结构:

Project(name)
 └── Filter(age > 30)
      └── Relation(users)

该结构由上至下描述了数据流动的方向,每一层均可独立定义变换规则。Catalyst提供了一个通用的 TreeNode 基类,所有计划节点均继承于此,并支持 transform 方法用于递归遍历并应用优化规则。

函数式变换机制的优势

Catalyst采用 纯函数式的树变换机制 ,即每次规则应用都会返回一个新的树实例,而非修改原对象。这种方式虽然带来一定的内存开销,但极大提升了系统的线程安全性与调试便利性。更重要的是,它允许使用组合子(Combinators)方式组织多个优化规则,形成“规则批次”(Rule Batches),按预定顺序依次执行。

以下是一个简化的Scala代码片段,展示如何通过 transform 方法实现常量折叠(Constant Folding)优化:

object ConstantFolding extends Rule[Expression] {
  def apply(expr: Expression): Expression = expr transform {
    case Add(Literal(c1, _), Literal(c2, _)) =>
      Literal(c1 + c2)
    case Multiply(Literal(0, _), _) =>
      Literal(0)
    case Equal(a, b) if a == b =>
      Literal(true)
  }
}

代码逻辑逐行解读:

  • 第1行:定义一个名为 ConstantFolding 的对象,继承自 Rule[Expression] ,表示这是一个作用于表达式类型的优化规则。
  • 第2行:重写 apply 方法,接收输入表达式 expr ,返回优化后的表达式。
  • 第3行:调用 transform 方法,开始对表达式树进行递归遍历。
  • 第4–5行:当遇到两个字面量相加时,直接计算结果并替换为新的字面量节点。
  • 第6–7行:乘法中若有一个操作数为0,则结果恒为0,提前折叠。
  • 第8–9行:相同表达式比较是否相等,可直接判定为 true

这类规则可以被批量注册进优化器阶段,构成完整的优化流水线。由于每条规则都是无副作用的纯函数,因此易于测试、复用和并行调度。

树结构带来的扩展优势

得益于树形建模方式,Catalyst能够轻松支持多种数据源和复杂查询结构。例如,在处理嵌套类型(如Parquet中的struct/array字段)时,只需扩展相应的表达式节点类型即可;而对于外部数据源(如Elasticsearch、MongoDB),也可通过自定义 Relation 节点接入优化流程。

此外,Catalyst还利用树结构实现了 逻辑计划验证 机制。在分析阶段,系统会检查属性引用是否合法、类型是否一致、聚合函数使用是否合规等,确保进入优化阶段的计划是语义正确的。这一过程同样基于模式匹配完成,体现了函数式编程在结构化数据处理中的强大表达能力。

graph TD
    A[SQL文本] --> B(Lexer/Parser)
    B --> C{AST}
    C --> D[Catalyst Analyzer]
    D --> E[Unresolved Logical Plan]
    E --> F[Catalog Lookup]
    F --> G[Resolved Logical Plan]
    G --> H[Catalyst Optimizer]
    H --> I[Optimized Logical Plan]
    I --> J[Physical Planner]
    J --> K[Selected Physical Plan]
    K --> L[Execution]

流程图说明:

上述Mermaid图展示了从原始SQL到最终执行计划的完整路径。其中,Analyzer阶段负责将未解析的AST转换为带有元数据信息的逻辑计划;Optimizer则在其基础上施加一系列规则变换;最后由物理规划器选择最优执行方案。整个流程以树结构贯穿始终,各阶段之间通过不可变对象传递状态,保证了系统的清晰边界与高内聚性。

阶段 输入 输出 主要任务
解析(Parsing) SQL字符串 抽象语法树(AST) 词法与语法分析
分析(Analysis) AST + Catalog 已解析逻辑计划 属性绑定、类型推断
优化(Optimization) 逻辑计划 优化后逻辑计划 应用规则简化表达式
物理规划(Physical Planning) 逻辑计划 可执行物理计划 选择Join策略、排序方式等

表格说明:

此表概括了Catalyst优化器的主要处理阶段及其职责分工。值得注意的是,“Catalog”在此指代元数据目录,包含表名、列名、数据类型等信息,是实现属性解析的关键依赖。

综上所述,Catalyst通过将查询建模为不可变树结构,并运用函数式编程原则进行递归变换,构建了一个高度模块化、可维护性强的优化框架。这种设计不仅降低了新增功能的技术门槛,也为后续引入成本模型奠定了良好基础。

2.1.2 规则与成本结合的优化策略设计

尽管基于规则的优化(Rule-Based Optimization, RBO)在早期版本中占据主导地位,但随着用户对查询性能要求的提高,单纯依赖启发式规则已难以应对复杂的现实工作负载。为此,Spark 2.2起逐步引入了 基于成本的优化 (Cost-Based Optimization, CBO),并与原有RBO机制形成互补,共同指导物理计划的选择。

规则驱动优化的特点与局限

RBO主要依赖预设的启发式规则来简化逻辑计划,典型示例如下:

  • 谓词下推 (Predicate Pushdown):将过滤条件尽可能靠近数据源执行,减少中间数据传输。
  • 投影裁剪 (Column Pruning):仅读取查询所需字段,避免加载冗余列。
  • 常量折叠 布尔简化 :提前计算静态表达式,降低运行时开销。
  • 子查询去关联化 (Unnest Subqueries):将相关子查询转换为连接操作,提升执行效率。

这些规则通常不依赖统计数据,执行速度快,适用于大多数常见场景。然而,它们无法回答诸如“哪个Join算法更优?”、“是否应该广播小表?”等问题,而这正是CBO的用武之地。

成本模型的核心要素

CBO通过估算不同执行路径的“代价”(通常以I/O、CPU、网络通信等资源消耗为单位),选择总成本最低的物理计划。在Catalyst中,成本评估主要依赖以下几类统计信息:

  • 行数估计 (Row Count)
  • 数据大小 (Data Size in Bytes)
  • 列基数 (Column Cardinality / Distinct Value Count)
  • 空值比例 (Null Ratio)
  • 最小/最大值范围 (Min/Max Values)

这些信息可通过执行 ANALYZE TABLE ... COMPUTE STATISTICS 命令收集,并存储在Hive Metastore或内置Catalog中。

假设存在两个表 orders (大表,1亿行)与 customers (小表,1万行),执行如下查询:

SELECT o.order_id, c.name 
FROM orders o JOIN customers c ON o.cust_id = c.id;

若未启用CBO,优化器可能默认采用Sort-Merge Join;但若开启CBO并发现 customers 表体积远小于阈值(默认 spark.sql.autoBroadcastJoinThreshold=10MB ),则会自动选择 Broadcast Hash Join ,显著减少Shuffle开销。

规则与成本的协同工作机制

Catalyst并非完全取代RBO,而是将其作为第一道防线,在逻辑计划层面尽可能消除冗余操作。随后,在物理规划阶段引入CBO进行细粒度决策。具体流程如下:

  1. 所有候选物理计划由 PhysicalPlanningStrategy 生成;
  2. 每个计划的成本由 CostEvaluator 计算;
  3. 最低成本的计划被选为最终执行方案。

此过程可通过配置参数控制,例如:

-- 启用CBO
spark.sql.cbo.enabled=true
spark.sql.cbo.joinReorder.enabled=true  # 允许重排多表连接顺序

下面是一段模拟成本评估的伪代码:

class CostBasedPlanner {
  def pickBestPlan(logicalPlan: LogicalPlan): PhysicalPlan = {
    val candidates = strategies.flatMap(_.apply(logicalPlan))
    candidates.minBy(costModel.estimate)
  }
}

object costModel {
  def estimate(plan: PhysicalPlan): Double = {
    val cpuCost = plan.expressions.map(_.complexity).sum
    val ioCost = plan.children.map(child => stats(child).sizeInBytes).sum
    val shuffleCost = if (plan.isShuffling) networkBandwidth * ioCost else 0
    cpuCost * 0.1 + ioCost * 0.5 + shuffleCost
  }
}

代码逻辑逐行解读:

  • pickBestPlan 方法接收逻辑计划,通过多个策略生成候选物理计划集合;
  • 使用 minBy 找出成本最低者;
  • costModel.estimate 综合考虑CPU、I/O和Shuffle三项成本;
  • 表达式复杂度、数据大小和是否触发Shuffle均为关键因子;
  • 权重系数可根据集群资源配置动态调整。
实际影响与调优建议

在实践中,混合优化策略显著提升了复杂查询的执行效率。例如,在TPC-DS基准测试中,启用CBO后部分查询性能提升可达 3倍以上 。但同时也带来额外元数据管理负担,需定期更新统计信息以避免误判。

推荐做法包括:
- 对频繁查询的大表定期执行 ANALYZE TABLE ... COMPUTE STATISTICS
- 利用 DESCRIBE TABLE EXTENDED 查看当前统计信息状态;
- 设置合理的广播阈值,防止OOM异常;
- 结合 EXPLAIN 命令观察实际选用的Join策略。

总之,Catalyst通过整合规则与成本两种优化范式,既保留了RBO的快速响应能力,又增强了对大规模数据集的适应性,构成了现代分布式SQL引擎的重要基石。

3. DAG调度系统与任务执行效率优化

Spark 作为分布式计算框架,其核心优势之一在于将用户编写的高阶操作(如 map filter join 等)自动转换为可并行执行的底层任务,并通过高效的调度机制实现资源最优利用。这一过程的关键组件是 DAGScheduler TaskScheduler ,它们共同构成了 Spark 的调度层级体系。在大规模数据处理场景中,作业性能瓶颈往往不在于计算本身,而在于任务调度不合理、并行度不足或数据本地性差等问题。因此,深入理解 DAG 调度系统的内部工作机制,掌握任务执行效率的调优方法,对于提升整体集群吞吐量和响应速度至关重要。

本章节将系统性地解析 Spark 的 DAG 调度流程,从阶段划分原理到任务分发策略,再到实际作业中的延迟根因分析,层层递进。我们将结合理论模型、源码逻辑、配置参数以及可视化工具(如 Spark UI),帮助读者构建完整的调度认知体系,并提供可落地的优化方案。

3.1 DAGScheduler的工作机制

DAGScheduler 是 Spark 调度层的核心模块,负责将用户的 RDD 操作图转化为可执行的 Stage 阶段序列,并管理这些阶段之间的依赖关系和提交顺序。它运行在 Driver 进程中,监听来自应用程序的操作请求,在 Action 触发时启动调度流程。整个机制基于有向无环图(Directed Acyclic Graph, DAG)建模,这也是“DAG”名称的由来。

3.1.1 阶段划分与宽窄依赖理论分析

在 Spark 中,每个 Job 对应一个 Action 操作(如 count() collect() saveAsTextFile() )。当 Job 被提交后,DAGScheduler 会根据该 Job 所涉及的 RDD 血统(Lineage)信息构建出一个 DAG 图,然后依据 窄依赖(Narrow Dependency) 宽依赖(Wide Dependency) 的边界进行阶段(Stage)划分。

  • 窄依赖 :父 RDD 的每个分区最多被子 RDD 的一个分区所使用,例如 map() filter() union()
  • 宽依赖 :父 RDD 的一个分区可能被多个子 RDD 分区所引用,通常发生在 Shuffle 操作中,如 groupByKey() reduceByKey() join()
graph TD
    A[TextFile] -->|map| B[(MappedRDD)]
    B -->|filter| C[(FilteredRDD)]
    C -->|reduceByKey| D[(ShuffledRDD)]
    D -->|map| E[(ResultRDD)]
    E -->|collect| F[Action]

    style A fill:#f9f,stroke:#333
    style F fill:#f96,stroke:#333

    subgraph Stages
        S1[Stage 0: Map -> Filter]
        S2[Stage 1: reduceByKey (Shuffle)]
        S3[Stage 2: Final Map]
    end

    S1 --> S2 --> S3

上图展示了典型的三阶段划分过程:
- Stage 0 包含从读取文件到 filter 的所有窄依赖操作;
- reduceByKey 引入了 Shuffle,形成宽依赖,成为阶段划分点;
- 后续的 map 操作无法与前一阶段合并,因此单独构成 Stage 2。

阶段划分算法逻辑

DAGScheduler 使用反向遍历的方式从最后一个 RDD 开始向上回溯:

private def createShuffleMapStages(shuffleDep: ShuffleDependency[_, _, _]): Unit = {
  getOrCreateParentStages(shuffleDep.rdd, shuffleDep.rdd.partitions.length)
  // 创建新的 ShuffleMapStage
}

逐行解读:
1. createShuffleMapStages 方法接收一个 ShuffleDependency 类型的依赖对象,表示存在 Shuffle 的边。
2. getOrCreateParentStages 递归查找上游 RDD 是否已经属于某个 Stage,若否,则创建新 Stage。
3. 最终生成的 Stage 分为两类:
- ShuffleMapStage :用于执行 Shuffle 写操作,输出中间结果供下游消费;
- ResultStage :最终执行 Action 的阶段,直接产生结果返回给 Driver。

参数说明与影响
参数 默认值 说明
spark.default.parallelism 根据集群核数决定 控制默认分区数,影响 Stage 中 Task 数量
spark.sql.shuffle.partitions 200 SQL 查询中 Join/Agg 的默认分区数
spark.stage.maxConsecutiveAttempts 4 单个 Stage 最大重试次数

⚠️ 注意:过多的小分区会导致 Task 调度开销上升;过少的大分区则容易造成内存溢出或长尾任务。

实际案例:Join 导致频繁 Stage 切分

假设我们有两个大表进行多次 join 操作:

val joined = tableA.join(tableB).join(tableC).count()

每次 join 都会触发一次 Shuffle,从而生成一个新的 Stage。如果未开启 Adaptive Query Execution(AQE),就会出现多个串行 Shuffle 阶段,显著增加总执行时间。

解决方案包括:
- 启用 AQE( spark.sql.adaptive.enabled=true ),允许运行时合并小 Stage;
- 使用广播 Join( broadcast() )减少 Shuffle;
- 调整 shuffle.partitions 以匹配数据规模。

3.1.2 Shuffle边界识别与任务分组策略

Shuffle 边界不仅是逻辑上的阶段切分点,更是物理执行层面的任务分组依据。DAGScheduler 在识别到宽依赖时,必须确保前一阶段的所有 Task 完成并将中间数据写入磁盘,才能启动下一阶段的 Task。这种“阻塞式推进”机制保障了数据一致性,但也带来了潜在的性能瓶颈。

Shuffle 边界的判定规则
操作类型 是否引入 Shuffle 原因
map , flatMap , filter 窄依赖,一对一映射
groupByKey , reduceByKey 需要跨节点聚合
join (非广播) 键值需重新分布
repartition(n) 显式重分区
coalesce(n, false) 否(若 n < 当前分区数) 只合并分区,不 Shuffle

可以通过以下代码判断某操作是否会触发 Shuffle:

def isShuffleOperation(rdd: RDD[_]): Boolean = {
  rdd.dependencies.exists(_.isInstanceOf[ShuffleDependency[_, _, _]])
}

逻辑分析:
1. rdd.dependencies 返回当前 RDD 的所有父依赖;
2. 使用 isInstanceOf[ShuffleDependency] 判断是否存在 Shuffle 类型依赖;
3. 若存在,说明该操作将引发网络传输和磁盘 I/O。

💡 提示:可通过 .toDebugString 查看血统链路:

scala println(rdd.toDebugString)

输出示例:
(8) MapPartitionsRDD[3] at map at Example.scala:10 [] | CoalescedRDD[2] at coalesce at Example.scala:9 [] | ShuffledRDD[1] at reduceByKey at Example.scala:8 [] | ParallelCollectionRDD[0] at parallelize at Example.scala:7 []

任务分组策略:Stage 内部的 Task 并行化

一旦确定 Stage 划分,DAGScheduler 将为每个 Stage 创建一组 Task(数量等于该 Stage 输入 RDD 的分区数),并将这些 Task 封装成 TaskSet 提交给 TaskScheduler。

reduceByKey 为例:

Stage ID 类型 输入分区数 Task 数 输出目标
0 ShuffleMapStage 100 100 写入本地块管理器
1 ResultStage 100 100 收集结果至 Driver

每个 Task 的职责明确:
- ShuffleMapTask :执行到 Shuffle 写为止,输出位于 BlockManager 的 shuffle_X_Y_Z.data 文件;
- ResultTask :拉取远程 Shuffle 数据,完成最终计算。

表格:不同操作对应的 Stage 结构对比
操作序列 DAG 结构 Stage 数 是否含 Shuffle
sc.textFile(...).map(...).count() Lineage: File → Map → Count 1
rdd.map(...).reduceByKey(...).collect() File → Map → [Shuffle] → Reduce → Collect 2
df.join(df2).groupBy(...).agg(...) (SQL) 多次 Shuffle ≥3
rdd.repartition(100).map(...) Repartition → Map 2
优化建议:减少不必要的 Shuffle

常见手段包括:
- 使用 combineByKeyWithClassTag 替代多次 groupByKey + mapValues
- 合理设置 spark.sql.autoBroadcastJoinThreshold (默认 10MB),让小表自动广播;
- 使用 bucketBy + sortWithinBuckets 预分区避免运行时 Shuffle。

3.2 TaskScheduler与资源分配逻辑

TaskScheduler 是 DAGScheduler 的下层调度器,负责将 TaskSet 中的任务分配给集群中的可用 Executor 执行。它对接底层资源管理器(如 YARN、Standalone 或 Kubernetes),实现任务的具体调度与失败恢复。

3.2.1 FIFO与FAIR调度模式对比

Spark 支持两种主要的调度池(Pool)模式:FIFO(先进先出)和 FAIR(公平调度),通过 spark.scheduler.mode 参数控制。

调度模式配置对照表
配置项 FIFO 模式 FAIR 模式
spark.scheduler.mode FIFO (默认) FAIR
任务排队方式 单队列,按提交顺序执行 多队列,支持权重与抢占
适用场景 批处理作业为主 多租户、交互式查询混合环境
并发执行能力 仅当前 Job 全部完成才释放资源 可并发执行多个 Job,动态共享 CPU
示例:启用 FAIR 调度

首先创建 fairscheduler.xml 配置文件:

<?xml version="1.0"?>
<allocations>
  <pool name="production">
    <schedulingMode>FAIR</schedulingMode>
    <weight>3</weight>
    <minShare>2</minShare>
  </pool>
  <pool name="development">
    <schedulingMode>FAIR</schedulingMode>
    <weight>1</weight>
    <minShare>1</minShare>
  </pool>
</allocations>

然后在 spark-submit 中指定:

--conf spark.scheduler.mode=FAIR \
--conf spark.scheduler.allocation.file=/path/to/fairscheduler.xml

参数解释:
- <weight> :决定资源分配比例,默认为 1,生产池权重为 3,意味着优先获得更多 CPU 时间片;
- <minShare> :最小保证资源量(以 Task 数计),即使其他队列空闲也不会剥夺;
- <schedulingMode> :可选 FAIR 或 FIFO。

性能对比实验
测试场景 FIFO 表现 FAIR 表现
单个大作业 ✅ 快速占满资源,执行快 ❌ 被拆分调度,略有延迟
多个小查询并发 ❌ 后面的查询长时间等待 ✅ 响应迅速,平均延迟低
混合负载(OLAP + 实时) ❌ 实时查询卡顿 ✅ 实时优先级高,体验好

🔍 推荐实践:在 BI 报表平台等多用户环境中启用 FAIR 模式,提升服务 SLA。

3.2.2 任务推测执行与失败重试机制

在异构集群中,某些节点可能因硬件老化、GC 停顿或网络拥塞导致 Task 执行缓慢(即“长尾问题”)。Spark 提供了 推测执行(Speculative Execution) 功能,自动检测慢 Task 并启动副本并行运行,取最先完成的结果。

开启推测执行的配置参数
参数 默认值 作用范围
spark.speculation false 全局开关
spark.speculation.interval 100ms 检查周期
spark.speculation.multiplier 3 慢于平均时间 X 倍视为慢 Task
spark.speculation.quantile 0.75 至少 75% 的 Task 完成才开始推测
示例配置:
--conf spark.speculation=true \
--conf spark.speculation.multiplier=2 \
--conf spark.speculation.quantile=0.8
源码级逻辑分析(简化版)
def checkSpeculatableTasks(): Unit = {
  val tasks = runningTasks.filter(_.runtime > median * multiplier)
  if (tasks.nonEmpty && completionRate > quantile) {
    tasks.foreach(launchTaskCopy(_))
  }
}

逐行解释:
1. runningTasks 获取当前正在运行的所有 Task;
2. 计算中位运行时间 median
3. 筛选出运行时间超过 median * multiplier 的 Task;
4. 判断已完成 Task 比例是否达到 quantile
5. 若满足条件,为每个慢 Task 启动副本( launchTaskCopy );

⚠️ 注意事项:
- 不适用于有副作用的操作(如 saveAsTextFile ),可能导致重复写入;
- 在 SSD 存储环境下效果更明显,HDD 容易误判;
- 建议配合监控指标(如 Task Time 分布)动态调整参数。

失败重试机制

当 Task 因异常退出时,Spark 默认会重试最多 4 次 (由 spark.task.maxFailures 控制):

--conf spark.task.maxFailures=6

若所有尝试均失败,则整个 Stage 失败,进而导致 Job 失败。此时可通过日志定位具体错误原因(如 OOM、序列化失败等)。

3.3 任务并行度设置与数据本地性优化

任务并行度与数据本地性是影响 Spark 作业性能的两大关键因素。合理的并行度可以最大化利用集群资源,而良好的本地性则能显著减少网络传输开销。

3.3.1 分区数与Executor核心数匹配原则

理想情况下,每个 Executor 应同时运行多个 Task,充分利用多核优势。基本原则如下:

总 Task 数 ≈ 总 CPU Core 数 × 2 ~ 3

例如:
- 集群有 10 个 Executor,每个 4 核 → 总共 40 核;
- 推荐设置分区数为 80~120。

常见设置方式
数据源 设置方法
RDD rdd.repartition(n) coalesce(n)
DataFrame df.repartition(n)
读取 HDFS 文件 文件块数 × spark.default.parallelism
自动计算推荐分区数的函数:
def recommendPartitions(executorNum: Int, coresPerExecutor: Int): Int = {
  val totalCores = executorNum * coresPerExecutor
  math.max(totalCores * 2, 200) // 至少 200 分区防止单 Task 过大
}

参数说明:
- executorNum : 实际启动的 Executor 数量;
- coresPerExecutor : 每个 Executor 分配的 vCore 数( spark.executor.cores );
- 返回值作为 repartition() 的输入。

实测案例:并行度过低导致资源浪费

某作业原始设置:
- 分区数 = 20
- Executor 数 = 10,每 Executor 4 core

结果:同一时刻只有 20 个 Task 运行,其余 20 个核心空闲,CPU 利用率不足 50%。

优化后:
- df.repartition(100)
- CPU 利用率提升至 90%以上,总耗时下降 60%。

3.3.2 数据本地性等级(NODE_LOCAL、PROCESS_LOCAL)控制

Spark 在调度 Task 时优先选择数据所在节点执行,以减少网络传输。本地性级别分为:

级别 说明 性能排序
PROCESS_LOCAL 数据与 Executor 在同一 JVM 最优
NODE_LOCAL 数据与 Executor 在同一节点 次优
NO_PREF 无偏好,任意调度 一般
RACK_LOCAL 同机架不同节点 较差
ANY 跨机架拉取数据 最差
调度等待策略

Spark 支持对不同本地性级别设置等待时间:

参数 默认值 含义
spark.locality.wait.process 3s 等待 PROCESS_LOCAL
spark.locality.wait.node 3s 等待 NODE_LOCAL
spark.locality.wait 3s 统一等待时间(覆盖前两者)
高效配置示例:
--conf spark.locality.wait=6s \
--conf spark.locality.wait.node=5s \
--conf spark.locality.wait.process=1s

📌 建议:在数据集中存储于 HDFS 的场景下,适当延长 node 等待时间,提高本地性命中率。

监控本地性统计

通过 Spark UI 的 Stage Details 页面查看:

Localty Level Tasks
PROCESS_LOCAL 150
NODE_LOCAL 30
ANY 2

ANY 比例过高(>10%),应检查:
- HDFS 块位置信息是否同步;
- Executor 是否动态分配导致位置偏移;
- 是否启用了 spark.scheduler.locationAwareTasks (默认 true)。

3.4 实战:大规模作业的任务延迟根因分析

在真实生产环境中,作业变慢的原因复杂多样。本节通过一个典型故障排查案例,演示如何使用 Spark UI 和关键配置调优解决调度瓶颈。

3.4.1 使用Spark UI定位调度瓶颈

假设某 ETL 作业耗时从 10min 上升至 40min,初步怀疑为调度问题。

步骤一:进入 Spark UI(http://driver:4040)

导航至 Executors Tab:
- 发现部分 Executor 内存使用接近 100%,且频繁 Full GC;
- Task Deserialization Time 达到 5s,远高于正常值(<100ms);

→ 推测:序列化开销大 + 内存不足。

步骤二:查看 Stage 页面

发现 Stage 5 出现大量红色 Task(失败或重试),且 Duration 分布极不均匀:
- 最快 Task:2s
- 最慢 Task:38s

→ 存在严重长尾问题。

步骤三:点击单个慢 Task 查看详情
  • Scheduler Delay : 12s → 表明 Task 提交后等待调度时间过长;
  • Fetch Wait Time : 8s → 说明 Shuffle 拉取阻塞;
  • Localty Level : ANY → 数据不在本地;

结论: 资源竞争激烈 + 数据本地性差 + Shuffle 热点

3.4.2 调整配置参数提升任务吞吐量

针对上述问题,实施以下优化措施:

优化清单
问题 解决方案 配置项
调度延迟高 增加并行度 spark.sql.shuffle.partitions=400
内存不足 增大堆外内存 spark.executor.memoryOverhead=1024
长尾 Task 开启推测执行 spark.speculation=true
本地性差 延长等待时间 spark.locality.wait=6s
Shuffle 拉取慢 调整缓冲区 spark.reducer.maxSizeInFlight=96m
最终提交命令示例:
spark-submit \
  --class com.example.ETLJob \
  --master yarn \
  --deploy-mode cluster \
  --num-executors 20 \
  --executor-cores 4 \
  --executor-memory 8g \
  --conf spark.sql.shuffle.partitions=400 \
  --conf spark.speculation=true \
  --conf spark.locality.wait=6s \
  --conf spark.reducer.maxSizeInFlight=96m \
  --conf spark.executor.memoryOverhead=1024 \
  etl-job.jar
效果验证
指标 优化前 优化后
总耗时 40min 12min
CPU 利用率 45% 88%
失败 Task 数 23 0
平均 Task Time 30s 8s

✅ 成功将作业性能提升 3.3 倍,具备稳定服务能力。

4. Spark与Hadoop 2.7集成体系构建

Apache Spark 作为当前主流的大数据处理框架,其在生产环境中的广泛应用离不开与 Hadoop 生态系统的深度整合。尤其在企业级部署中,Hadoop 2.7 版本因其稳定性、成熟的 YARN 资源管理能力以及 HDFS 的高可用性,成为许多组织的首选底层基础设施。将 Spark 部署于 Hadoop 2.7 环境下,不仅能够复用现有的集群资源和存储系统,还能实现统一的资源调度、安全认证与运维监控。本章深入探讨 Spark 如何与 Hadoop 2.7 构建稳定高效的集成体系,涵盖从数据读取、资源调度到安全控制的全链路技术细节。

在整个集成架构中,HDFS 扮演着统一数据源的角色,为 Spark 提供持久化、可扩展的数据访问能力;YARN 则承担了分布式资源管理职责,确保 Spark 应用能够在多租户环境下公平、高效地运行;而 Kerberos 安全机制则保障了跨组件通信的身份可信性。这种“计算-资源-存储-安全”四位一体的架构设计,构成了现代大数据平台的核心骨架。理解并掌握这一集成体系,是构建企业级数据中台或分析平台的关键前提。

更进一步地,随着业务规模的增长,静态资源配置已难以满足动态负载需求,因此引入基于 YARN 的动态资源分配机制显得尤为重要。该机制允许 Spark 根据实际工作负载自动扩缩 Executor 实例,从而提升资源利用率并降低运营成本。同时,在安全集群中提交作业时,必须正确处理 Kerberos 认证流程与 Token 传递逻辑,否则会导致作业失败或权限异常。这些高级特性虽增加了部署复杂度,但却是保障系统稳定性和合规性的必要手段。

为了帮助读者建立完整的认知路径,本章将依次解析 HDFS 数据接入优化、YARN 上的部署模式选择、动态资源调度策略以及安全认证集成等关键环节,并结合真实场景提供配置示例与代码实践。通过深入剖析每个子模块的技术原理与交互机制,展示如何在复杂环境中实现 Spark 与 Hadoop 的无缝协同。

4.1 HDFS作为统一数据源的技术实现

HDFS(Hadoop Distributed File System)是 Hadoop 生态中最核心的分布式文件系统,具备高容错性、高吞吐量和横向扩展能力,非常适合用于大规模批处理任务的数据存储。当 Spark 运行在 Hadoop 集群之上时,HDFS 自然成为其默认且最常用的统一数据源。Spark 通过内置的 Hadoop InputFormat 接口与 HDFS 进行交互,实现了对多种文件格式的原生支持,并能根据文件分块策略智能划分 RDD 分区,最大化并行处理效率。

4.1.1 文件分块机制与RDD分区映射关系

HDFS 在写入大文件时会将其切分为固定大小的数据块,默认块大小为 128MB(Hadoop 2.7 中可通过 dfs.blocksize 参数调整)。每个数据块会被复制到多个 DataNode 上以保证可靠性,通常副本数为 3。这种分块机制不仅是 HDFS 实现分布式存储的基础,也为上层计算引擎如 Spark 提供了天然的并行处理单元。

Spark 在读取 HDFS 文件时,会依据文件的 block 分布信息来创建 RDD 的分区。具体而言,每一个 HDFS block 对应一个 RDD partition,这意味着 Spark 能够实现“数据本地性”(data locality),即任务尽量在存储该 block 的节点上执行,避免网络传输开销。这一映射关系由 TextInputFormat 或其他 Hadoop InputFormat 子类驱动,最终由 HadoopRDD 类完成分区构建。

下面是一个典型的 Spark 读取 HDFS 文本文件的代码示例:

val conf = new SparkConf().setAppName("HDFSRDDExample")
val sc = new SparkContext(conf)

// 从HDFS读取文本文件
val textFile = sc.textFile("hdfs://namenode:9000/user/data/input.txt")

// 统计每行单词数量
val wordCount = textFile.flatMap(_.split("\\s+"))
                       .map(word => (word, 1))
                       .reduceByKey(_ + _)

wordCount.saveAsTextFile("hdfs://namenode:9000/user/output/wordcount")
代码逻辑逐行解读:
  • 第1行:初始化 Spark 配置对象,设置应用名称。
  • 第2行:创建 SparkContext,连接到集群。
  • 第5行:调用 textFile() 方法从 HDFS 路径加载数据。此方法内部使用 TextInputFormat 将文件按 block 切分,并生成对应数量的 RDD 分区。
  • 第8–10行:执行标准的 WordCount 操作,包括分词、映射和聚合。
  • 第12行:将结果保存回 HDFS,输出路径也会被划分为多个 part 文件,每个对应一个 task 输出。
参数说明:
  • hdfs://namenode:9000 是 NameNode 的地址和端口,需根据实际集群配置修改。
  • sc.textFile() 支持第二个参数指定最小分区数,例如 sc.textFile(path, 4) 可强制至少生成 4 个分区,即使文件小于 4 个 block。
  • 若文件压缩(如 .gz ),Spark 会退化为单分区读取,因为 gzip 不支持 split,影响并行度。

下表展示了不同文件格式与是否可分割之间的关系及其对 Spark 分区的影响:

文件格式 是否可分割 InputFormat 实现 对 Spark 分区影响
TextFile TextInputFormat 每个 block 对应一个 partition
SequenceFile SequenceFileInputFormat 支持记录级同步点,block 可分割
Avro AvroInputFormat 块内包含 sync marker,支持 split
Parquet ParquetInputFormat 列式存储,row group 级别可并行读取
Gzip (.gz) TextInputFormat(受限) 整个文件作为一个 partition,串行读取

注意 :不可分割的压缩格式会严重限制并行度,建议在大数据场景中优先使用支持 split 的压缩方式,如 Snappy + SequenceFile 或 LZO。

此外,Spark 还提供了 coalesce() repartition() 方法来手动调整分区数。但在理想情况下,应尽可能依赖 HDFS block 到 RDD partition 的自然映射,以保持最优的数据本地性。

graph TD
    A[HDFS File] --> B{Splitable?}
    B -->|Yes| C[Block 1 → Partition 1]
    B -->|Yes| D[Block 2 → Partition 2]
    B -->|No| E[Entire File → Single Partition]
    C --> F[Task Executed on Node with Block]
    D --> F
    E --> G[Task May Run Remotely]

该流程图清晰表达了 HDFS 文件分块如何影响 Spark 的任务调度决策。只有当文件可分割时,才能实现真正的分布式并行读取。

4.1.2 SequenceFile、Parquet等格式读写优化

除了文本文件外,企业在生产环境中更多采用二进制序列化格式进行高效存储与交换。其中,SequenceFile 和 Parquet 是两种典型代表,分别适用于不同的使用场景。

SequenceFile 是 Hadoop 原生的键值对存储格式,适合中间数据传输或 MapReduce 输出。它支持三种记录类型: Uncompressed , Record Compressed , Block Compressed 。推荐使用 Block Compressed 模式,因为它会对多个记录打包后统一压缩,显著提高压缩率和 I/O 性能。

以下是在 Spark 中读写 SequenceFile 的 Scala 示例:

// 写入 SequenceFile
val data = List(("key1", "value1"), ("key2", "value2"))
val rdd = sc.parallelize(data)
rdd.saveAsSequenceFile("hdfs://namenode:9000/user/output/seqfile")

// 读取 SequenceFile
val seqRdd = sc.sequenceFile[String, String](
  "hdfs://namenode:9000/user/output/seqfile",
  classOf[String],
  classOf[String]
)
代码解释:
  • saveAsSequenceFile() 使用 Writable 类型序列化键值对,默认使用 Text BytesWritable
  • sequenceFile[K,V]() 需显式指定键和值的 Writable 类型,Spark 会通过反射实例化解码器。
  • 若需自定义类型,需实现 Writable 接口并注册 Kryo 序列化器。

相比之下, Parquet 是一种面向分析的列式存储格式,特别适合 OLAP 查询场景。它支持谓词下推(predicate pushdown)、列裁剪(column pruning)和高效的压缩编码(如 RLE、Dictionary),极大减少了磁盘 I/O 和内存占用。

使用 Spark SQL 读写 Parquet 文件极为简便:

import org.apache.spark.sql.functions._

val df = spark.read.parquet("hdfs://namenode:9000/user/data/users.parquet")

// 执行列裁剪和过滤下推
val filteredDf = df.select("name", "age")
                   .filter(col("age") > 30)

filteredDf.write.mode("overwrite").parquet("hdfs://namenode:9000/user/output/adults")
执行逻辑分析:
  • read.parquet() 触发元数据读取(来自 _metadata 或 footer),仅加载所需 schema。
  • select() 实现列裁剪,后续操作只读取 name 和 age 列。
  • filter() 条件会被下推至文件扫描阶段,利用 Parquet 行组(Row Group)的统计信息跳过不满足条件的 block。
  • 写入时可通过配置启用 Snappy 压缩: spark.sql.parquet.compression.codec=snappy
优化项 配置参数 推荐值 作用说明
Parquet 压缩算法 spark.sql.parquet.compression.codec snappy 平衡压缩比与速度
并行读取阈值 spark.sql.files.maxPartitionBytes 134217728 (128MB) 控制每个 task 处理的最大数据量
推断分区字段 spark.sql.sources.partitionInference.enabled true 自动识别 /date=2024-01-01/ 类型路径为分区列
缓存 Parquet 元数据 spark.sql.parquet.enable.summary-metadata true 加速后续查询的 schema 获取

上述配置应在 spark-defaults.conf 中预先设定,以确保所有作业继承最佳实践。

此外,对于频繁访问的小表,可考虑启用 Parquet 的 Bloom Filter Index Page 功能(需使用 Delta Lake 或 Apache Iceberg 扩展),实现更快的点查性能。

综上所述,合理选择文件格式并配合相应的读写优化策略,不仅能提升 Spark 作业的整体执行效率,也能有效降低 HDFS 的 I/O 压力。特别是在混合负载环境中,应根据数据特征(结构化程度、访问模式、更新频率)灵活选用合适的存储格式。

4.2 YARN上的Spark应用部署模式

YARN(Yet Another Resource Negotiator)是 Hadoop 2.x 引入的核心资源管理框架,负责集群中 CPU、内存等资源的统一调度与隔离。Spark 可以作为 YARN 上的一个客户端应用程序运行,借助其强大的资源分配能力实现多租户共享、优先级调度和故障恢复。然而,Spark on YARN 提供了两种主要的部署模式:Client 模式和 Cluster 模式,二者在驱动程序(Driver)的位置、生命周期管理和容错能力方面存在本质差异,直接影响作业的稳定性与运维体验。

4.2.1 Client与Cluster模式的区别与选型

Client 模式 下,Spark Driver 运行在提交作业的客户端机器上,而 Executor 则由 YARN 分配在各个 NodeManager 上。这意味着客户端必须持续保持在线状态,直到整个应用结束。一旦客户端断开连接(如 SSH 会话中断),Driver 进程将终止,导致整个 Spark 应用失败。

而在 Cluster 模式 中,Driver 本身也被封装进一个 YARN Container,并由 ResourceManager 调度到某个 NodeManager 上运行。此时客户端仅负责提交应用后即可退出,后续的 Driver 生命周期由 YARN 监控和维护,具备更强的容错能力。

对比维度 Client 模式 Cluster 模式
Driver 运行位置 提交节点(本地 JVM) 集群内部某个 NodeManager
客户端依赖 必须保持连接 提交后可断开
日志查看方式 直接打印到本地终端 需通过 yarn logs -applicationId <id> 查看
网络延迟敏感性 高(Driver 与 Executor 跨网络通信) 低(Driver 与 Executor 同处内网)
适用场景 开发调试、交互式 Shell(spark-shell) 生产环境批处理任务

选择哪种模式取决于具体的使用场景。例如,在开发阶段使用 spark-shell spark-submit 时,常采用 Client 模式以便实时观察日志输出;而在生产环境中调度定时 ETL 任务时,则强烈推荐使用 Cluster 模式,以避免因客户端宕机导致作业中断。

以下是使用 spark-submit 提交作业的命令示例:

# Client 模式提交
spark-submit \
  --master yarn \
  --deploy-mode client \
  --name MySparkApp \
  --class com.example.WordCount \
  /path/to/myapp.jar \
  hdfs://namenode:9000/input \
  hdfs://namenode:9000/output

# Cluster 模式提交
spark-submit \
  --master yarn \
  --deploy-mode cluster \
  --name MySparkApp \
  --class com.example.WordCount \
  --conf spark.yarn.submit.waitAppCompletion=false \
  /path/to/myapp.jar \
  hdfs://namenode:9000/input \
  hdfs://namenode:9000/output
参数说明:
  • --master yarn :指定运行在 YARN 上。
  • --deploy-mode :设置为 client cluster
  • --conf spark.yarn.submit.waitAppCompletion=false :在 Cluster 模式下,客户端不等待应用完成,立即返回。
  • 若未设置该参数, spark-submit 将阻塞直至应用结束,不适合后台调度。

4.2.2 Container资源申请与JVM堆外内存配置

当 Spark 运行在 YARN 上时,所有资源请求都通过 YARN Container 实现。每个 Executor 实际上运行在一个 YARN Container 内,Container 的资源规格由 Spark 配置决定。若配置不当,可能导致资源浪费或容器被 YARN 拒绝。

关键资源配置参数如下:

参数名 默认值 说明
spark.executor.memory 1g Executor JVM 堆内存大小
spark.executor.memoryOverhead max(384, 0.1 * heap) 堆外内存(用于 Netty、Python、字符串等)
spark.driver.memory 1g Driver 堆内存
spark.driver.memoryOverhead 同上 Driver 堆外内存
spark.executor.cores 1 每个 Executor 使用的 CPU 核心数
spark.executor.instances 动态 Executor 实例总数
yarn.nodemanager.resource.memory-mb 8192 单个 NodeManager 总可用内存(需提前规划)

假设我们希望启动一个拥有 4 个 Executor、每个使用 4 核 CPU 和 8GB 堆内存的应用,则配置如下:

spark-submit \
  --master yarn \
  --deploy-mode cluster \
  --num-executors 4 \
  --executor-cores 4 \
  --executor-memory 8g \
  --driver-memory 4g \
  --conf spark.executor.memoryOverhead=2048 \
  --conf spark.driver.memoryOverhead=1024 \
  --class com.example.AnalyticsJob \
  /apps/spark/jobs/analytics.jar
内存计算逻辑:
  • 每个 Executor 请求的总内存 = executor-memory + memoryOverhead
  • 即:8192 MB + 2048 MB = 10240 MB
  • YARN Container 最小单位通常为 1024 MB,因此实际分配为 11 * 1024 ≈ 11264 MB(向上取整)
  • 若 NodeManager 总内存为 16GB(16384 MB),最多容纳 1 个这样的 Executor(剩余不足)

为了避免资源碎片化,建议遵循以下原则:
1. spark.executor.memory 不超过 NodeManager 内存的 1/2;
2. spark.executor.cores 设置为虚拟核数的整除因子(如 4 核机器设为 2 或 4);
3. memoryOverhead 至少设置为 executor-memory 的 10%,复杂任务可增至 20%。

pie
    title YARN Container 内存组成(10GB Executor)
    “JVM Heap” : 8192
    “Off-heap (MemoryOverhead)” : 2048

此外,还需注意 YARN 的队列容量限制。若目标队列(如 default prod )已满,应用将处于 ACCEPTED 状态无法启动。可通过 yarn queue -status <queue> 检查队列使用情况。

综上,合理的资源规划是保障 Spark on YARN 稳定运行的前提。不仅要关注 Spark 自身配置,还需与 YARN 的全局资源策略协同设计,避免出现“申请不到 container”或“频繁 OOM”的问题。

5. 大规模数据处理平台全流程实战

5.1 从零搭建Spark+Hadoop一体化平台

在构建企业级大数据处理系统时,搭建一个稳定、可扩展的 Spark + Hadoop 一体化平台是实现高效批流一体计算的基础。本节将从环境准备到核心配置文件调优,完整演示平台搭建流程。

5.1.1 环境准备与节点角色规划

假设我们部署一个包含 5 台物理节点的集群(均为 CentOS 7.x,JDK 1.8+),其角色分配如下:

节点名称 IP 地址 角色职责
node1 192.168.1.101 NameNode, ResourceManager, Master
node2 192.168.1.102 SecondaryNameNode, HistoryServer
node3 192.168.1.103 DataNode, NodeManager, Worker
node4 192.168.1.104 DataNode, NodeManager, Worker
node5 192.168.1.105 DataNode, NodeManager, Worker

部署前需完成以下准备工作:
- 所有节点关闭防火墙并配置免密 SSH 登录;
- 安装 JDK 并设置 JAVA_HOME
- 配置 NTP 时间同步服务,避免 Kerberos 认证失败;
- 下载 Apache Hadoop 2.7.7 和 Spark 3.1.1 编译版本(支持 Hadoop 2.7);

解压后建议统一路径:

tar -zxf hadoop-2.7.7.tar.gz -C /opt/
tar -zxf spark-3.1.1-bin-hadoop2.7.tgz -C /opt/

5.1.2 配置文件详解(spark-defaults.conf, yarn-site.xml)

Hadoop 核心配置片段( yarn-site.xml
<configuration>
    <!-- 启用容量调度器 -->
    <property>
        <name>yarn.resourcemanager.scheduler.class</name>
        <value>org.apache.hadoop.yarn.server.resourcemanager.scheduler.capacity.CapacityScheduler</value>
    </property>

    <!-- Container内存限制 -->
    <property>
        <name>yarn.nodemanager.resource.memory-mb</name>
        <value>16384</value>
    </property>
    <property>
        <name>yarn.scheduler.maximum-allocation-mb</name>
        <value>8192</value>
    </property>

    <!-- JVM 堆外内存配置 -->
    <property>
        <name>yarn.nodemanager.vmem-pmem-ratio</name>
        <value>2.1</value>
    </property>

    <!-- 日志聚合 -->
    <property>
        <name>yarn.log-aggregation-enable</name>
        <value>true</value>
    </property>
</configuration>

参数说明
- yarn.nodemanager.resource.memory-mb :每节点可用内存总量(16GB);
- vmem-pmem-ratio :虚拟内存与物理内存比率,防止因堆外内存超限被 Kill;

Spark 配置( spark-defaults.conf
spark.master                     yarn
spark.submit.deployMode          cluster
spark.executor.instances         3
spark.executor.cores             4
spark.executor.memory            6g
spark.driver.memory              4g
spark.serializer                 org.apache.spark.serializer.KryoSerializer
spark.sql.adaptive.enabled       true
spark.dynamicAllocation.enabled  true
spark.shuffle.service.enabled    true

关键优化点解析
- 使用 Kryo 序列化提升网络传输效率;
- 开启自适应查询执行(AQE)以动态合并小分区;
- 启用动态资源分配,配合 YARN 实现弹性扩缩容;

部署完成后,可通过以下命令验证集群状态:

# 启动 HDFS 和 YARN
$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh

# 提交测试作业
$SPARK_HOME/bin/spark-submit \
  --class org.apache.spark.examples.SparkPi \
  --master yarn \
  --num-executors 2 \
  $SPARK_HOME/examples/jars/spark-examples*.jar 10

作业提交后可在 YARN Web UI(http://node1:8088)查看运行状态,确认容器正常启动且无 OOM 错误。

该平台架构为后续 ETL、流处理和机器学习任务提供了统一运行时基础,实现了存储与计算分离的设计理念。

本文还有配套的精品资源,点击获取 menu-r.4af5f7ec.gif

简介:Spark 3.1.1 是 Apache Spark 的一个重要版本,结合 Hadoop 2.7 的分布式存储与资源管理能力,构建了高性能、高可靠的大数据处理环境。本文深入解析 Spark 在 Hadoop 平台上的核心特性,涵盖 SQL 性能优化、内存计算、实时流处理、机器学习支持及容错机制等关键内容。通过 Spark 与 YARN、HDFS 的深度集成,实现资源高效调度与大规模数据处理,适用于数据分析、实时计算和 AI 工作负载。本项目可作为企业级大数据平台搭建的参考实践。


本文还有配套的精品资源,点击获取
menu-r.4af5f7ec.gif

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐