基于python如何使用pyspark?
·
大家好,我是jobleap.cn的小九。今天学习了pyspark的基础用法。
下面以销售数据分析为典型场景,提供一个PySpark快速入门示例,涵盖环境准备、核心操作和分析流程。
一、环境准备
PySpark依赖Java环境,需先确保:
- 安装Java 8/11(推荐8),并配置
JAVA_HOME环境变量 - 安装PySpark:
pip install pyspark
二、核心示例:销售数据分析
场景目标:处理一份销售记录数据,完成「数据查看→过滤高销售额记录→按地区聚合总销售额」的分析流程。
完整代码
# 1. 导入必要模块
from pyspark.sql import SparkSession
from pyspark.sql.functions import sum
# 2. 初始化SparkSession(PySpark的入口点)
spark = SparkSession.builder \
.appName("SalesAnalysisQuickstart") # 应用名称
.master("local[*]") # 本地模式,利用所有可用CPU核心
.getOrCreate()
# 3. 准备数据(模拟销售记录:日期、产品、销售额、地区)
sales_data = [
("2023-01-01", "手机", 3500, "华东"),
("2023-01-02", "电脑", 6000, "华北"),
("2023-01-03", "手机", 4200, "华东"),
("2023-01-04", "平板", 2800, "华南"),
("2023-01-05", "电脑", 5500, "华北"),
("2023-01-06", "平板", 3000, "华南")
]
# 定义列名
columns = ["date", "product", "sales", "region"]
# 4. 创建DataFrame(分布式数据表,类似Pandas的DataFrame)
sales_df = spark.createDataFrame(data=sales_data, schema=columns)
# 5. 数据探索
print("=== 数据结构 ===")
sales_df.printSchema() # 查看字段类型
print("\n=== 原始数据 ===")
sales_df.show() # 展示前20行数据(分布式数据的本地预览)
# 6. 数据处理:过滤高销售额记录(销售额>4000)
high_sales_df = sales_df.filter(sales_df.sales > 4000)
print("\n=== 高销售额记录(>4000) ===")
high_sales_df.show()
# 7. 数据分析:按地区聚合总销售额
region_total_sales = sales_df.groupBy("region") \
.agg(sum("sales").alias("total_sales")) \
.orderBy("total_sales", ascending=False) # 按总销售额降序
print("\n=== 各地区总销售额 ===")
region_total_sales.show()
# 8. 停止SparkSession(释放资源)
spark.stop()
三、代码解释
-
SparkSession初始化:
SparkSession是PySpark的核心入口,负责创建分布式数据集(DataFrame)和执行操作。local[*]表示在本地模式运行,适合开发测试。 -
DataFrame创建:通过
createDataFrame将本地数据转换为分布式DataFrame,schema定义字段名和类型(自动推断类型,也可手动指定)。 -
数据探索:
printSchema():查看字段结构(类似数据库表结构)show():本地预览数据(默认前20行)
-
数据处理:
filter():按条件筛选记录(类似SQL的WHERE)- 支持链式调用,语法接近Pandas
-
数据分析:
groupBy():按字段分组agg():聚合操作(如sum计算总和)orderBy():排序结果
四、运行结果
执行后会输出:
- 数据结构(字段名和类型)
- 原始销售记录
- 销售额>4000的记录(电脑和部分手机)
- 各地区总销售额(华东:7700,华北:11500,华南:5800)
五、扩展说明
- 实际场景中,数据通常来自文件(CSV/Parquet)或数据库,可通过
spark.read.csv("path")或spark.read.jdbc(...)读取。 - PySpark的优势是处理大规模数据(远超单机内存),自动分布式计算。
- 更多操作可参考官方文档:PySpark SQL
更多推荐
所有评论(0)