小菜鸟

java菜鸟号正在起航

spark历史服务器:追踪任务全生命周期

Spark 跑任务的时候,4040 端口能看到执行详情——DAG、Stage、Task、Shuffle 数据量。

但任务一结束,4040 就关了。你还没来得及看,窗口没了。

历史服务器解决的就是这个问题:把任务日志持久化存下来,跑完也能看。

  • 应用跑的时候 → 把事件日志写到指定目录(HDFS 或本地)
  • 应用跑完之后 → 历史服务器从那个目录读日志,生成 Web UI
  • 你在 18080 端口能看到所有历史任务

配置分两步:应用写日志 + 服务读日志

这两步必须配同一个目录——应用往那里写,历史服务器从那里读。

Spark 应用 → 写日志到 /spark-history → 历史服务器 → 从 /spark-history 读 → 展示在 18080

第一步:让 Spark 应用写日志(spark-defaults.conf

cd $SPARK_HOME/conf
cp spark-defaults.conf.template spark-defaults.conf

在文件里加两行:

spark.eventLog.enabled  true
spark.eventLog.dir      hdfs://localhost:9000/spark-history
参数 说明
spark.eventLog.enabled 开日志收集,默认是关的
spark.eventLog.dir 日志写到哪(HDFS 或本地路径)

路径支持两种:

# HDFS 路径(生产推荐,分布式存储)
spark.eventLog.dir  hdfs://localhost:9000/spark-history

# 本地路径(单机测试用)
spark.eventLog.dir  file:///usr/local/spark/logs/history

提前创建目录并给权限:

# HDFS
hadoop fs -mkdir -p /spark-history
hadoop fs -chmod 777 /spark-history

# 本地
mkdir -p /usr/local/spark/logs/history
chmod 777 /usr/local/spark/logs/history

第二步:启动历史服务器读日志(spark-env.sh

阅读全文 »

Spring 事务失效场景深度解析

Spring 事务基于 AOP 动态代理实现,事务失效的本质是:事务相关的 AOP 增强逻辑未被触发。本文从 “代理机制” 切入,逐一拆解 6 种常见失效场景的底层原因,并给出可落地的解决方案,尤其重点讲解 “this 调用失效” 这一难点。

事务失效的核心前提:理解 Spring 事务的 AOP 原理

在分析失效场景前,必须先明确 Spring 事务的执行流程(基于 CGLIB 代理,最常见场景):

  1. 代理类生成:Spring 为标注 @Transactional 的目标 Bean(如 UserService)生成 CGLIB 代理类(继承自目标 Bean);
  2. 外部调用触发增强:当外部代码调用代理类的方法时,代理类会先执行 事务增强逻辑(开启事务),再调用目标 Bean 的原始方法;
  3. 目标方法执行:目标 Bean 执行业务逻辑,若正常完成则代理类触发事务提交,若抛出异常则触发回滚;
  4. 内部调用绕过增强:若目标 Bean 内部通过 this 调用自身方法,this 指向目标 Bean 本身(非代理类),会直接执行原始方法,跳过代理类的增强逻辑 → 事务失效。

6 种事务失效场景与底层解析

1. 场景 1:事务方法非 public 修饰

底层原因

Spring AOP 动态代理(JDK/CGLIB)仅对 public 方法生效:

  • JDK 代理:基于接口实现,仅代理接口中的 public 方法;
  • CGLIB 代理:基于继承生成子类,非 public 方法(private/protected/default)无法被重写,代理类无法插入增强逻辑。

Spring 源码佐证(AbstractFallbackTransactionAttributeSource):

protected TransactionAttribute computeTransactionAttribute(Method method, @Nullable Class<?> targetClass) {
    // 非 public 方法直接返回 null,无事务属性
    if (Modifier.isPublic(method.getModifiers()) == false) {
        return null;
    }
    // ... 后续逻辑
}

若方法非 public,Spring 会忽略其 @Transactional 注解,不生成事务增强。

阅读全文 »

Spark 环境配置:本地调试用 Local,小集群用 Standalone,生产上 YARN

Spark 有三种运行模式,按复杂度从低到高:

模式 适用场景 特点
Local 开发调试、单元测试 单机跑,不用配置集群,最简单
Standalone 小规模集群、快速部署 Spark 自带的集群管理,不需要 Hadoop
YARN 生产环境、已有 Hadoop 集群 跟 HDFS 无缝集成,统一资源管理

Spark 有三种运行模式,按复杂度从低到高:

模式 适用场景 特点
Local 开发调试、单元测试 单机跑,不用配置集群,最简单
Standalone 小规模集群、快速部署 Spark 自带的集群管理,不需要 Hadoop
YARN 生产环境、已有 Hadoop 集群 跟 HDFS 无缝集成,统一资源管理

前置准备:Java + Spark 安装包

不管哪种模式,都需要先装好 Java 和 Spark。

Java 版本: JDK 8 或 11(推荐 8,兼容性最好)

下载 Spark:

Spark 安装包分两种:

包名 特点 适合谁
spark-x.x.x-bin-hadoopx.x 自带 Hadoop 依赖 新手、快速测试
spark-x.x.x-bin-without-hadoop 不带 Hadoop,需要手动关联 已有 Hadoop 环境
# 解压
tar -zxvf spark-3.1.1-bin-without-hadoop.tgz -C /usr/local/
mv /usr/local/spark-3.1.1-bin-without-hadoop /usr/local/spark

环境变量(~/.bash_profile):

export SPARK_HOME=/usr/local/spark
export PATH=$SPARK_HOME/bin:$PATH

生效:source ~/.bash_profile

如果用 without-hadoop 版本,还需要关联 Hadoop 类路径:

cd $SPARK_HOME/conf
cp spark-env.sh.template spark-env.sh

spark-env.sh 里加一行:

export SPARK_DIST_CLASSPATH=$(hadoop classpath)

Local 模式:什么都不用配,直接跑

Local 模式是所有进程在单个 JVM 里跑,没有集群,没有分布式调度,用来验证 Spark 能不能用。

启动 Spark Shell:

spark-shell

成功的话会看到:

阅读全文 »

Spark 入门实战:WordCount 程序全解析,从代码到执行

WordCount 是大数据的 Hello World——统计文本里每个单词出现几次。

Spark 版本的 WordCount 核心就四步:

  1. 读文件 → RDD[String]
  2. 拆单词 → flatMap
  3. 变键值对 → map(word => (word, 1))
  4. 聚合求和 → reduceByKey( + )

这四步写完,Spark 帮你分布式并行执行

环境准备:版本匹配与依赖配置

Spark 对 Scala 版本有严格依赖,错误的版本组合会导致兼容性问题(如类找不到、方法异常),需提前确认版本对应关系。

版本选择原则

  • Spark 2.x 主要支持 Scala 2.11、2.12;
  • Spark 3.x 主要支持 Scala 2.12、2.13(但部分早期 3.x 版本对 2.13 支持不完善);
  • 推荐组合:Spark 3.1.1 + Scala 2.12.x(稳定性好,生态支持完善)。

maven版本对应

Maven 依赖配置

pom.xml 中添加 Spark Core 和 Scala 依赖:

<!-- Scala 核心库 -->  
<dependency>  
  <groupId>org.scala-lang</groupId>  
  <artifactId>scala-library</artifactId>  
  <version>2.12.13</version>  <!-- 与 Spark 版本匹配 -->  
</dependency>  

<!-- Spark 核心依赖 -->  
<dependency>  
  <groupId>org.apache.spark</groupId>  
  <artifactId>spark-core_2.12</artifactId>  <!-- _2.12 表示适配 Scala 2.12 -->  
  <version>3.1.1</version>  
</dependency>

注意:spark-core artifactId 中的 _2.12 必须与 Scala 版本一致,否则会出现 ClassNotFoundException

代码:四步完成计数

import org.apache.spark.{SparkConf, SparkContext}

object WordCount {
  def main(args: Array[String]): Unit = {
    // 1. 配置
    val conf = new SparkConf()
      .setMaster("local[*]")   // 本地模式,用所有 CPU 核心
      .setAppName("WordCount")

    // 2. 创建上下文(Spark 的"入口")
    val sc = new SparkContext(conf)

    // 3. 读文件 → RDD[String](每行一个元素)
    val lines = sc.textFile("src/main/resources/wordcount.txt")

    // 4. 拆单词 → flatMap(一行 → 多个单词)
    val words = lines.flatMap(_.split(" "))

    // 5. 变键值对 → map(每个单词 → (单词, 1))
    val wordAndOne = words.map(word => (word, 1))

    // 6. 聚合 → reduceByKey(相同 key 的 value 相加)
    val wordCounts = wordAndOne.reduceByKey(_ + _)

    // 7. 触发执行 + 打印结果
    wordCounts.collect().foreach(println)

    // 8. 关闭
    sc.stop()
  }
}

输入文件内容(wordcount.txt):

Hadoop Hive Spark
Hadoop HDFS
Hive Spark
Zookeeper Kafka Flume

输出:

(Hadoop, 2)
(Hive, 2)
(Spark, 2)
(HDFS, 1)
(Zookeeper, 1)
(Kafka, 1)
(Flume, 1)

关键概念:三步走

1. SparkConf + SparkContext = 启动入口

组件 做什么
SparkConf 配置运行模式、应用名称、资源参数
SparkContext 连接 Spark 集群,创建 RDD,调度任务

SparkContext 是所有 Spark 程序的”大门”——有了它才能操作 RDD。

2. RDD:分布式数据集合

textFile 读进来的 lines 是一个 RDD。它不是一个本地集合,是分布在多台机器上的数据分片。

RDD 的核心特性:

  • 不可变:创建后不能改,要改就生成新 RDD
  • 分区:数据分散在多台机器上
  • 惰性:不遇到 Action 不会真的算

3. Transformation vs Action

类型 特点 例子
Transformation 惰性执行,只记录”要做什么”,返回新 RDD flatMapmapreduceByKey
Action 触发真正执行,返回结果或写外部存储 collectcountsaveAsTextFile

代码里的执行时机:

lines.flatMap(...)     ← 还没算,只是记下了
  .map(...)            ← 还没算,只是记下了
  .reduceByKey(...)    ← 还没算,只是记下了
  .collect()           ← 触发执行!所有前面的操作开始跑
阅读全文 »

Spark 与 Hadoop:不是替代,是”存储 + 计算”的分工

Hadoop 和 Spark 经常被放在一起比较,但它们的关系不是”谁取代谁”——更像基础设施和上层应用

层次 Hadoop 提供 Spark 提供
存储 HDFS(分布式文件系统) 依赖外部存储,不自己管存
计算 MapReduce(磁盘级批处理) 内存计算引擎
资源调度 YARN 可以跑在 YARN 上

最常用的组合:HDFS 存数据 + YARN 管资源 + Spark 算数据

先看定位:一个管存,一个管算

Hadoop Spark
核心定位 分布式存储 + 批量计算 通用计算引擎
存储能力 自带 HDFS 依赖 HDFS/S3/Kafka
计算模型 磁盘级批处理(MapReduce) 内存级计算(RDD/DAG)
资源管理 自带 YARN 可独立,也可跑在 YARN/K8s
适用场景 离线批处理、冷数据存储 批处理、流处理、SQL、机器学习

Hadoop 的核心是”存”——它是一个存储底座。Spark 的核心是”算”——它是一个计算引擎。

计算层对比:MapReduce vs Spark

这是两者最直接的差异。

对比项 MapReduce Spark
中间结果 写磁盘(HDFS) 优先内存,不够才写磁盘
多阶段任务 每阶段写一次磁盘 DAG 串联,内存流转
任务启动 每个任务起 JVM(秒级) 线程池复用(毫秒级)
排序策略 强制排序 按需排序
迭代计算 每次读磁盘,效率低 内存缓存,效率高
延迟 分钟-小时级 秒-分钟级

一个直观的例子:

阅读全文 »
0%