PySpark 入门:分布式数据处理与 SQL 查询
·
PySpark 入门:分布式数据处理与 SQL 查询
PySpark 是 Apache Spark 的 Python API,支持分布式数据计算。其核心优势在于:
- 分布式处理:自动将数据分区并行处理,适合海量数据集
- 统一引擎:批处理、流计算、机器学习均可实现
- 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 $$ 。
更多推荐



所有评论(0)