PySpark 入门:分布式数据处理与 SQL 查询

PySpark 是 Apache Spark 的 Python API,支持分布式数据计算。其核心优势在于:

  1. 分布式处理:自动将数据分区并行处理,适合海量数据集
  2. 统一引擎:批处理、流计算、机器学习均可实现
  3. SQL 集成:通过 DataFrame API 直接执行 SQL 查询

1. 环境初始化

创建 SparkSession(入口对象):

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("FirstApp") \  # 应用名称
    .getOrCreate()         # 创建或复用会话


2. 分布式数据处理

创建 DataFrame(核心数据结构):

# 从列表创建分布式数据集
data = [("Alice", 34), ("Bob", 45), ("Cathy", 29)]
columns = ["Name", "Age"]
df = spark.createDataFrame(data, columns)

# 查看数据分布
df.show()

输出:

+-----+---+
| Name|Age|
+-----+---+
|Alice| 34|
|  Bob| 45|
|Cathy| 29|
+-----+---+

转换操作示例(惰性执行):

# 过滤年龄 > 30 的记录
df_filtered = df.filter(df.Age > 30)

# 添加新列
from pyspark.sql.functions import lit
df_new = df.withColumn("Country", lit("USA"))


3. SQL 查询集成

注册临时视图

df.createOrReplaceTempView("people")  # 创建SQL可访问的视图

执行 SQL 查询

result = spark.sql("""
    SELECT Name, Age 
    FROM people 
    WHERE Age BETWEEN 25 AND 40
""")
result.show()

输出:

+-----+---+
| Name|Age|
+-----+---+
|Alice| 34|
|Cathy| 29|
+-----+---+


4. 性能优化技巧
  • 分区策略:通过 repartition(4) 控制数据分区数
  • 缓存机制:对复用数据集调用 df.cache() 减少重复计算
  • 谓词下推:SQL 查询自动优化为最小化数据传输

5. 完整示例:数据分析流程
# 1. 读取CSV文件(分布式加载)
df = spark.read.csv("data.csv", header=True, inferSchema=True)

# 2. 注册SQL视图
df.createOrReplaceTempView("sales")

# 3. 执行聚合查询
spark.sql("""
    SELECT product, AVG(price) as avg_price 
    FROM sales 
    GROUP BY product
""").show()

# 4. 关闭会话
spark.stop()

关键概念

  • RDD:底层分布式数据集(Resilient Distributed Dataset)
  • DataFrame:结构化数据抽象,支持 SQL 查询
  • Catalyst 优化器:自动优化查询执行计划
    通过 PySpark 可实现 TB 级数据处理,速度比传统单机工具快 $$ \text{10-100} \times $$ 。
Logo

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

更多推荐