更多 Spark DataFrame 内容

PySpark 入门

Benjamin Schmidt

Data Engineer

从多种数据源创建 DataFrame

  • CSV 文件:常见的结构化、分隔数据
  • JSON 文件:半结构化、层级格式
  • Parquet 文件:为存储与查询优化,常用于数据工程
  • 示例:
    spark.read.csv("path/to/file.csv")
    
  • 示例:
    spark.read.json("path/to/file.json")
    
  • 示例:
    spark.read.parquet("path/to/file.parquet")
    
1 https://spark.apache.org/docs/latest/api/python/reference/pyspark.pandas/api/pyspark.pandas.read_csv
PySpark 入门

模式推断与手动定义

  • inferSchema=True 时,Spark 可从数据推断模式

  • 为更好控制可手动定义模式——适用于固定结构数据

大规模模式

PySpark 入门

PySpark DataFrame 的数据类型

  • IntegerType:整数
    • 例如:13478-1890456
  • LongType:更大的整数
    • 例如:8 字节有符号数,922334775806
  • FloatType 与 DoubleType:用于小数的浮点数
    • 例如:3.14159
  • StringType:文本/字符串数据
    • 例如:"This is an example of a string."
  • ...
PySpark 入门

PySpark DataFrame 数据类型语法

# 以类形式导入必要类型
from pyspark.sql.types import (StructType,
                            StructField, IntegerType,
                            StringType, ArrayType)

# 构建模式
schema = StructType([
    StructField("id", IntegerType(), True),
    StructField("name", StringType(), True),
    StructField("scores", ArrayType(IntegerType()), True)
])

# 设置模式
df = spark.createDataFrame(data, schema=schema)
PySpark 入门

DataFrame 操作:选择与筛选

  • 使用 .select() 选择特定列
  • 使用 .filter().where() 按条件筛选行
  • 使用 .sort() 按多列排序
# 仅选择并显示 name 与 age 列
df.select("name", "age").show()
# 筛选 age > 30
df.filter(df["age"] > 30).show()
# 使用 where 精确匹配值
df.where(df["age"] == 30).show()
# 使用 sort 按 age 降序
df.sort("age", ascending=False).show()
PySpark 入门

排序与删除缺失值

  • 使用 .sort().orderBy() 排序
  • 使用 na.drop() 删除含空值的行
# 按 age 列排序
df.sort("age", ascending=False).show()

# 删除缺失值
df.na.drop().show()

PySpark 入门

速查表

  • spark.read_json(): 从 JSON 加载数据
  • spark.read.schema(): 显式定义模式
  • .na.drop(): 丢弃含缺失值的行
  • .select(), .filter(), .sort(), .orderBy(): 基础数据操作函数
PySpark 入门

Let's practice!

PySpark 入门

Preparing Video For Download...