大家好,我是jobleap.cn的小九。今天学习了pyspark的基础用法。
下面以销售数据分析为典型场景,提供一个PySpark快速入门示例,涵盖环境准备、核心操作和分析流程。

一、环境准备

PySpark依赖Java环境,需先确保:

  1. 安装Java 8/11(推荐8),并配置JAVA_HOME环境变量
  2. 安装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()

三、代码解释

  1. SparkSession初始化SparkSession是PySpark的核心入口,负责创建分布式数据集(DataFrame)和执行操作。local[*]表示在本地模式运行,适合开发测试。

  2. DataFrame创建:通过createDataFrame将本地数据转换为分布式DataFrame,schema定义字段名和类型(自动推断类型,也可手动指定)。

  3. 数据探索

    • printSchema():查看字段结构(类似数据库表结构)
    • show():本地预览数据(默认前20行)
  4. 数据处理

    • filter():按条件筛选记录(类似SQL的WHERE)
    • 支持链式调用,语法接近Pandas
  5. 数据分析

    • groupBy():按字段分组
    • agg():聚合操作(如sum计算总和)
    • orderBy():排序结果

四、运行结果

执行后会输出:

  • 数据结构(字段名和类型)
  • 原始销售记录
  • 销售额>4000的记录(电脑和部分手机)
  • 各地区总销售额(华东:7700,华北:11500,华南:5800)

五、扩展说明

  • 实际场景中,数据通常来自文件(CSV/Parquet)或数据库,可通过spark.read.csv("path")spark.read.jdbc(...)读取。
  • PySpark的优势是处理大规模数据(远超单机内存),自动分布式计算。
  • 更多操作可参考官方文档:PySpark SQL

更多推荐