PySpark 入门
Ben Schmidt
Data Engineer
.na.drop() 删除包含空值的行# 删除包含任一空值的行 df_cleaned = df.na.drop()# 过滤掉空值 df_cleaned = df.where(col("columnName").isNotNull())
.na.fill({"column": value) 将空值替换为指定值# 将 age 列中的空值填充为 0
df_filled = df.na.fill({"age": 0})
.withColumn() 基于计算或现有列添加新列# 新增列 'age_plus_5'
df = df.withColumn("age_plus_5", df["age"] + 5)
withColumnRenamed() 重命名列# 将 'age' 列重命名为 'years'
df = df.withColumnRenamed("age", "years")
drop() 删除不需要的列# 删除 'department' 列
df = df.drop("department")
.filter() 按条件选择行# 筛选 salary 大于 50000 的行
filtered_df = df.filter(df["salary"] > 50000)
.groupBy() 与聚合函数(如 .sum(), .avg())汇总数据# 按部门分组并计算平均薪资
grouped_df = df.groupBy("department").avg("salary")
过滤
+------+---+-----------------+
|salary|age| occupation |
+------+---+-----------------+
| 60000| 45|Exec-managerial |
| 70000| 35|Prof-specialty |
+------+---+-----------------+
分组(GroupBy)
`
+----------+-----------+
|department|avg(salary)|
+----------+-----------+
| HR| 80000.0|
| IT| 70000.0|
+----------+-----------+
`
# 删除包含任一空值的行 df_cleaned = df.na.drop()# 在指定列上删除空值 df_cleaned = df.where(col("columnName").isNotNull())# 将 age 列中的空值填充为 0 df_filled = df.na.fill({"age": 0})
使用 .withColumn() 基于计算或现有列添加新列。语法:.withColumn("new_col_name", "original transformation")
# 新增列 'age_plus_5'
df = df.withColumn("age_plus_5", df["age"] + 5)
使用 withColumnRenamed() 重命名列
语法:withColumnRenamed(old column name,new column name`
# 将 'age' 列重命名为 'years'
df = df.withColumnRenamed("age", "years")
drop() 删除不需要的列
语法:.drop(column name)# 删除 'department' 列
df = df.drop("department")
# 筛选 salary 大于 50000 的行
filtered_df = df.filter(df["salary"] > 50000)
.groupBy() 与聚合函数(如 .sum(), .avg())汇总数据 # 按部门分组并计算平均薪资
grouped_df = df.groupBy("department").avg("salary")
PySpark 入门