Spark SQL 实战:WordCount 与商品销售分析
Spark SQL 入门项目,用 Scala 实现 WordCount 和商品销售总额分析,涵盖 DataFrame/SQL 两种 API。
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 优化原理和分区/缓存等优化手段。