PySpark 入门
Benjamin Schmidt
Data Engineer
spark.read.csv("path/to/file.csv")
spark.read.json("path/to/file.json")
spark.read.parquet("path/to/file.parquet")
在 inferSchema=True 时,Spark 可从数据推断模式
为更好控制可手动定义模式——适用于固定结构数据

IntegerType:整数1、3478、-18904569223347758063.14159"This is an example of a string."# 以类形式导入必要类型
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)
.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()
.sort() 或 .orderBy() 排序na.drop() 删除含空值的行# 按 age 列排序
df.sort("age", ascending=False).show()
# 删除缺失值
df.na.drop().show()
spark.read_json(): 从 JSON 加载数据spark.read.schema(): 显式定义模式.na.drop(): 丢弃含缺失值的行.select(), .filter(), .sort(), .orderBy(): 基础数据操作函数PySpark 入门