PySpark DataFrame 简介

PySpark 入门

Benjamin Schmidt

Data Engineer

关于 DataFrame

  • DataFrame:表格形式(行/列)
  • 支持类 SQL 操作
  • 类似 Pandas DataFrame 或 SQL 表
  • 结构化数据

Dataframes

PySpark 入门

从文件存储创建 DataFrame

# 从 CSV 创建 DataFrame
census_df = spark.read.csv('path/to/census.csv', header=True, inferSchema=True)
PySpark 入门

打印 DataFrame

# 显示前 5 行
census_df.show()


   age  education.num marital.status         occupation income
0   90              9        Widowed                  ?  <=50K
1   82              9        Widowed    Exec-managerial  <=50K
2   66             10        Widowed                  ?  <=50K
3   54              4       Divorced  Machine-op-inspct  <=50K
4   41             10      Separated     Prof-specialty  <=50K
PySpark 入门

打印 DataFrame 模式

# 显示 schema
census_df.printSchema()

输出: root |-- age: integer (nullable = true) |-- education.num: integer (nullable = true) |-- marital.status: string (nullable = true) |-- occupation: string (nullable = true) |-- income: string (nullable = true)
PySpark 入门

PySpark DataFrame 的基础分析

# .count() 返回 DataFrame 的总行数
row_count = census_df.count()
print(f'Number of rows: {row_count}')
# groupby() 可进行类 SQL 聚合
census_df.groupBy('gender').agg({'salary_usd': 'avg'}).show()

其他聚合函数包括:

  • sum()
  • min()
  • max()
PySpark 入门

PySpark 分析常用函数

  • .select():选择特定列
  • .filter():按条件筛选行
  • .groupBy():按一列或多列分组
  • .agg():对分组数据做聚合
PySpark 入门

常用函数示例

# 用 filter 和 select 缩小 DataFrame 范围
filtered_census_df = census_df.filter(df['age'] > 50).select('age', 'occupation')
filtered_census_df.show()

输出 +---+------------------+ |age| occupation | +---+------------------+ | 90| ?| | 82| Exec-managerial| | 66| ?| | 54| Machine-op-inspct| +---+------------------+
PySpark 入门

开始练习吧!

PySpark 入门

Preparing Video For Download...