2024-08-10



import scala.concurrent.Future
import scala.concurrent.ExecutionContext.Implicits.global
import scala.util.{Success, Failure}
import play.api.libs.ws._
import play.api.libs.json._
import scala.concurrent.duration._
 
// 假设以下方法用于获取亚马逊商品页面的HTML
def fetchProductPage(asin: String, proxy: Option[String] = None): Future[String] = {
  val url = s"http://www.amazon.com/dp/$asin/?tag=yourtag-20"
  val request = WS.url(url)
  proxy.foreach(request.withProxy)
  val response = request.get()
 
  response.map { res =>
    if (res.status == 200) {
      res.body
    } else {
      throw new Exception(s"Failed to fetch product page for ASIN: $asin, status: ${res.status}")
    }
  }
}
 
// 使用示例
val asin = "B01M8L5Z3Q" // 示例ASIN
val proxyOption = Some("http://user:password@proxyserver:port") // 代理服务器(如有需要)
 
val pageFuture = fetchProductPage(asin, proxyOption)
 
pageFuture.onComplete {
  case Success(html) => println(s"Success: $html")
  case Failure(e) => println(s"Failed: ${e.getMessage}")
}
 
// 等待响应,如果需要同步执行
import scala.concurrent.Await
Await.result(pageFuture, 30.seconds) match {
  case html => println(s"Success: $html")
}

这个代码示例展示了如何使用Scala和Play WS库来异步获取亚马逊商品页面的HTML内容。它使用Future来处理异步操作,并且可以通过可选的代理服务器参数来绕过反爬虫措施。这个例子简洁地展示了如何应对代理和反爬虫的挑战,同时保持代码的简洁性和可读性。

2024-08-07

Spark 经典demo 的 Scala 和 Java 实现

一、背景与问题

在大数据处理领域,Spark 是一个核心的分布式计算框架,其核心抽象 RDD(Resilient Distributed Dataset)和 DAG(Directed Acyclic Graph)调度模型是理解其运行机制的关键。本文将通过 Spark 的经典 demo,深入探讨其工作原理,并通过 Scala 和 Java 两种语言实现对比,分析其适用场景和注意事项。

Spark 的核心优势在于其内存计算能力,能够将中间结果缓存于内存中,大幅提高处理效率。然而,这种优势也伴随着一些限制,例如内存占用过高可能导致 OOM(Out Of Memory)错误,或者在处理小数据量时反而不如传统批处理工具(如 MapReduce)高效。

二、基本原理

1. RDD 的核心概念

RDD 是 Spark 的核心数据结构,具有以下特点:

  • 分布式性:数据被分割成多个分区(Partition),分布在集群的不同节点上。
  • 惰性求值:所有转换操作(Transformation)都是惰性的,直到遇到 Action 操作(如 count()save())才会实际执行。
  • 容错性:通过 lineage(血缘)记录数据的生成过程,当某一分区数据丢失时,可以重新计算。

2. DAG 调度模型

Spark 通过 DAG(有向无环图)调度器将任务划分为 Stage,每个 Stage 包含多个 Task。DAG 调度器会根据数据的分区位置和依赖关系,优化任务的执行顺序,最大化数据本地性(Data Locality)。

3. 核心操作分类

  • Transformation:惰性操作(如 mapfiltergroupBy
  • Action:触发计算(如 count()reduce()save()

三、环境准备

1. 环境要求

  • Spark 3.x(推荐 3.3.0)
  • Java 8 或 11
  • Scala 2.12 或 2.13(根据 Spark 版本选择)
  • IDE:IntelliJ IDEA 或 VS Code(推荐 Scala 插件)

2. 初始化 Spark 环境

# 创建项目目录
mkdir spark-demo && cd spark-demo

# 初始化 Maven 项目(Java 示例)
mvn archetype:generate -DarchetypeArtifactId=maven-archetype-quickstart -DgroupId=com.example -DartifactId=spark-demo -DinteractiveMode=false

# 初始化 Scala 项目(Scala 示例)
sbt new scala/scala-seed.g8

四、核心实现

1. Scala 实现:Word Count

示例代码

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

object WordCountScala {
  def main(args: Array[String]): Unit = {
    // 初始化 Spark 配置
    val conf = new SparkConf().setAppName("WordCountScala").setMaster("local[*]")
    val sc = new SparkContext(conf)

    // 读取文本文件(本地或 HDFS)
    val textRDD = sc.textFile("src/main/resources/input.txt")

    // 转换操作:拆分单词并统计
    val wordCounts = textRDD
      .flatMap(line => line.split("\\W+")) // 将每行拆分为单词
      .filter(word => word.nonEmpty)        // 过滤空字符串
      .map(word => (word, 1))               // 转换为 (word, 1)
      .reduceByKey(_ + _)                  // 按单词聚合

    // Action 操作:输出结果
    wordCounts.foreach(println)

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

关键代码解释

  • flatMap:将每行文本拆分为单词,返回一个 RDD[String]。
  • filter:去除空字符串(如标点符号),避免统计错误。
  • map:将每个单词转换为 (word, 1) 元组,为后续聚合做准备。
  • reduceByKey:在集群中按 key 聚合值,使用 + 操作符累加计数。
  • foreach:触发计算并输出结果。

2. Java 实现:Word Count

示例代码

import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.api.java.function.Function;
import org.apache.spark.sql.SparkConf;

public class WordCountJava {
    public static void main(String[] args) {
        // 初始化 Spark 配置
        SparkConf conf = new SparkConf().setAppName("WordCountJava").setMaster("local[*]");
        JavaSparkContext sc = new JavaSparkContext(conf);

        // 读取文本文件
        JavaRDD<String> textRDD = sc.textFile("src/main/resources/input.txt");

        // 转换操作:拆分单词并统计
        JavaRDD<String> wordsRDD = textRDD.flatMap(new Function<String, Iterable<String>>() {
            @Override
            public Iterable<String> call(String line) {
                return Arrays.asList(line.split("\\W+"));
            }
        });

        JavaRDD<Tuple2<String, Integer>> wordCountsRDD = wordsRDD.map(new Function<String, Tuple2<String, Integer>>() {
            @Override
            public Tuple2<String, Integer> call(String word) {
                return new Tuple2<>(word, 1);
            }
        }).reduceByKey((a, b) -> a + b);

        // Action 操作:输出结果
        wordCountsRDD.foreach(System.out::println);

        // 关闭 SparkContext
        sc.stop();
    }
}

关键代码解释

  • flatMap:使用 Function 接口实现单词拆分,返回 Iterable<String>
  • map:将单词转换为 (word, 1) 元组,使用 Tuple2 类。
  • reduceByKey:使用 lambda 表达式 (a, b) -> a + b 实现计数聚合。
  • foreach:触发计算并输出结果。

3. Scala vs Java 实现对比

特性Scala 实现Java 实现
语法简洁性更简洁,支持函数式编程需要显式定义类和接口
类型推断支持类型推断需要显式声明类型
可读性更易读,适合数据处理任务代码量较大,适合复杂逻辑
性能略优(编译器优化)相当(JIT 编译优化)
学习成本需掌握函数式编程概念传统面向对象编程更易上手

五、完整案例

案例:日志分析系统

需求

分析服务器日志,统计每个 IP 的访问次数,并找出访问量最高的前 10 个 IP。

实现步骤

  1. 读取日志文件(每行格式:IP - - [01/Jan/2023:12:34:56 +0800] "GET /index.html HTTP/1.1" 200 1234
  2. 提取 IP 地址
  3. 统计访问次数
  4. 排序并输出前 10 个结果

Scala 实现代码

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

object LogAnalysisScala {
  def main(args: Array[String]): Unit = {
    val conf = new SparkConf().setAppName("LogAnalysisScala").setMaster("local[*]")
    val sc = new SparkContext(conf)

    val logRDD = sc.textFile("src/main/resources/logs.txt")

    val ipCounts = logRDD
      .map(line => {
        // 提取 IP 地址(假设日志格式固定)
        val parts = line.split(" ")
        val ip = parts(0)
        (ip, 1)
      })
      .reduceByKey(_ + _)

    val top10 = ipCounts
      .sortBy(_._2, false)  // 按访问次数降序排序
      .take(10)

    top10.foreach(println)

    sc.stop()
  }
}

Java 实现代码

import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.api.java.function.Function;
import org.apache.spark.sql.SparkConf;

public class LogAnalysisJava {
    public static void main(String[] args) {
        SparkConf conf = new SparkConf().setAppName("LogAnalysisJava").setMaster("local[*]");
        JavaSparkContext sc = new JavaSparkContext(conf);

        JavaRDD<String> logRDD = sc.textFile("src/main/resources/logs.txt");

        JavaRDD<Tuple2<String, Integer>> ipCountsRDD = logRDD.map(new Function<String, Tuple2<String, Integer>>() {
            @Override
            public Tuple2<String, Integer> call(String line) {
                // 提取 IP 地址
                String[] parts = line.split(" ");
                String ip = parts[0];
                return new Tuple2<>(ip, 1);
            }
        }).reduceByKey((a, b) -> a + b);

        // 排序并取前 10
        JavaRDD<Tuple2<String, Integer>> top10 = ipCountsRDD
            .sortBy(new Function<Tuple2<String, Integer>, Double>() {
                @Override
                public Double call(Tuple2<String, Integer> tuple) {
                    return -tuple._2;  // 按访问次数降序
                }
            }).take(10);

        top10.forEach(System.out::println);

        sc.stop();
    }
}

六、源码解析

1. SparkContext 的初始化

val conf = new SparkConf().setAppName("WordCountScala").setMaster("local[*]")
val sc = new SparkContext(conf)
  • setMaster("local[*]"):在本地运行,使用所有 CPU 核心。
  • setAppName:设置应用名称,用于集群管理界面查看。

2. RDD 的转换操作

val wordCounts = textRDD
  .flatMap(line => line.split("\\W+"))
  .filter(word => word.nonEmpty)
  .map(word => (word, 1))
  .reduceByKey(_ + _)
  • flatMap:将每行拆分为单词,返回一个 RDD[String]。
  • filter:去除空字符串,避免统计错误。
  • map:将单词转换为 (word, 1) 元组。
  • reduceByKey:在集群中按 key 聚合值,使用 + 操作符累加。

3. Action 操作的触发

wordCounts.foreach(println)
  • foreach 是 Action 操作,触发 RDD 的计算,返回结果。

七、进阶使用

1. 使用 Spark SQL 进行结构化处理

import org.apache.spark.sql.SparkSession

object SQLExample {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder
      .appName("SQLExample")
      .master("local[*]")
      .getOrCreate()

    val df = spark.read.text("src/main/resources/input.txt")

    df.createOrReplaceTempView("words")

    val wordCounts = spark.sql("SELECT word, COUNT(*) as count FROM words GROUP BY word")
    wordCounts.show()
  }
}
  • Spark SQL 提供了更高级的接口,适合处理结构化数据。
  • 使用 SQL 查询可以提高代码可读性,但需要熟悉 SQL 语法。

2. 使用 DataFrame 和 Dataset 进行优化

val df = spark.read.text("src/main/resources/input.txt")
val wordCountsDF = df
  .withColumn("word", split(col("value"), "\\W+").getItem(0))
  .filter(col("word").isNotNull)
  .groupBy("word")
  .agg(count("*").alias("count"))
  • withColumn:添加新列,提取单词。
  • groupByagg:进行聚合操作,支持 SQL 语法。
  • 使用 DataFrame API 可以利用 Spark 的优化器进行查询计划优化。

八、性能与工程实践

1. 分区策略优化

  • 默认分区数:Spark 会根据集群配置自动计算分区数,但可能需要手动调整。
  • 自定义分区:使用 repartitioncoalesce 调整分区数。
val repartitioned = textRDD.repartition(10)
  • 分区数选择:通常设置为集群核心数的 2-3 倍,避免过多的小文件。

2. 持久化策略

  • 缓存策略:使用 cache()persist() 缓存中间结果,避免重复计算。
  • 存储级别:选择合适的存储级别(如 MEMORY_AND_DISK)。
val cachedRDD = textRDD.map(...).cache()

3. 并行度调整

  • 并行度:通过 setExecutorMemoryOverheadsetExecutorCores 调整执行器配置。
  • 任务数:通过 getNumPartitionsgetNumPartitions 控制任务数量。

4. 数据本地性优化

  • 数据本地性:Spark 优先将任务分配到数据所在的节点,减少网络传输。
  • 数据倾斜:使用 repartitionsalting 解决数据倾斜问题。

九、常见问题与踩坑

1. 数据倾斜问题

现象:某些分区的数据量远大于其他分区,导致任务执行时间不均。

解决方法

  • 使用 repartitioncoalesce 重新分区。
  • 使用 salting 技术,将数据分散到多个分区。
val saltedRDD = textRDD.map { line =>
  val salt = (line.hashCode % 100).toString
  (salt, line)
}.partitionBy(new RandomPartitioner(sc.getConf, 100))

2. 内存不足导致的 OOM 错误

现象:程序运行时内存不足,导致 JVM 崩溃。

解决方法

  • 增加堆内存:--driver-memory--executor-memory
  • 使用 persist(StorageLevel.MEMORY_AND_DISK) 将数据存储到磁盘。

3. 分区数过少导致性能下降

现象:分区数太少,导致任务并行度不足。

解决方法

  • 使用 repartition 增加分区数。
  • 调整 spark.sql.shuffle.partitions 配置。

4. 任务调度开销过大

现象:任务调度时间过长,影响整体性能。

解决方法

  • 使用 checkpoint 中断长链式依赖。
  • 启用 spark.locality.wait 调整数据本地性等待时间。

十、最佳实践

1. 合理选择存储级别

  • 内存优先:使用 MEMORY_ONLYMEMORY_AND_DISK 缓存中间结果。
  • 磁盘存储:对于大数据量,使用 DISK_ONLY 避免内存溢出。

2. 使用惰性求值优化计算

  • 避免在转换操作中提前触发计算,直到遇到 Action 操作。

3. 分区策略与数据量匹配

  • 小数据量使用默认分区,大数据量手动调整分区数。

4. 使用 Spark SQL 进行结构化处理

  • 对结构化数据使用 SQL 查询,提高可读性和性能。

5. 监控和调优

  • 使用 Spark UI 监控任务执行情况,调整配置参数。

十一、总结

Spark 是一个强大的分布式计算框架,其核心抽象 RDD 和 DAG 调度模型是其高效运行的关键。通过 Scala 和 Java 的实现对比,可以看出 Scala 在表达复杂逻辑时更加简洁,而 Java 更适合需要严格类型控制的场景。在实际项目中,Spark 适用于处理大规模数据、需要内存计算的场景,但在小数据量或需要低延迟的场景中需谨慎使用。通过合理调整分区策略、使用缓存和持久化策略,可以显著提升性能。同时,需要注意数据倾斜、内存不足等常见问题,通过优化配置和代码结构,确保 Spark 任务的稳定性和效率。