跳到正文
Joeplover
学习笔记·2025-10-24·约 2 分钟阅读

Spark SQL 实战:WordCount 与商品销售分析

Spark SQL 入门项目,用 Scala 实现 WordCount 和商品销售总额分析,涵盖 DataFrame/SQL 两种 API。

Spark 分析

Spark SQL 实战:WordCount 与商品销售分析

DataFrame

DataFrame 在 RDD 基础上增加了 Schema 信息——每一列的名字和类型。这看起来只是小改变,但带来的性能提升是巨大的。

# RDD 方式:Spark 不知道字段类型
rdd = sc.textFile("sales.csv")     .map(lambda line: line.split(","))     .filter(lambda fields: float(fields[2]) > 100)

# DataFrame 方式:Spark 知道字段类型和结构
df = spark.read.csv("sales.csv", header=True, inferSchema=True)
df.filter(df.amount > 100)

前者 Spark 只能序列化/反序列化整个对象,后者 Spark 知道只需要读取 amount 这一列(列裁剪)。

Catalyst 优化器

Spark SQL 的性能秘密在于 Catalyst 优化器。它会对查询计划做多种优化:

SELECT name, amount FROM sales WHERE amount > 100

列裁剪:只读 name 和 amount 两列,忽略其他列(减少 I/O)。

谓词下推:在读取数据时就把 amount > 100 的过滤条件应用到数据源(减少读取量)。

# 两种写法等价,Catalyst 优化后执行计划相同
df.createOrReplaceTempView("sales")
spark.sql("SELECT product, SUM(amount) as total FROM sales GROUP BY product ORDER BY total DESC")

# 等价的 DataFrame API
df.groupBy("product").agg(sum("amount").alias("total")).orderBy(desc("total"))

商品销售分析实例

from pyspark.sql import SparkSession
from pyspark.sql.functions import sum, count, avg, col

spark = SparkSession.builder.appName("SalesAnalysis").getOrCreate()

# 读取销售数据
sales = spark.read.csv("sales.csv", header=True, inferSchema=True)
products = spark.read.csv("products.csv", header=True, inferSchema=True)

# 热门商品 Top 10
top_products = (sales
    .groupBy("product_id")
    .agg(sum("amount").alias("total_sales"))
    .orderBy(col("total_sales").desc())
    .limit(10))

# 关联商品信息
result = (top_products
    .join(products, "product_id")
    .select("product_name", "total_sales", "category"))

# 按品类汇总
category_stats = (result
    .groupBy("category")
    .agg(
        sum("total_sales").alias("category_total"),
        count("*").alias("product_count")
    ))

category_stats.show()

学习建议

Spark SQL 的学习路径:DataFrame API → SQL → 性能调优。

DataFrame API 覆盖 90% 以上的数据分析场景。学会它,每天处理的数据分析需求基本都能搞定。当性能成为瓶颈时,再去深入研究 Catalyst 优化原理和分区/缓存等优化手段。