创建列
在本章中,您将学习如何使用 Spark 的 DataFrame 类所定义的方法来执行常见的数据操作。
先来看按列进行的操作。在 Spark 中,您可以使用 .withColumn() 方法完成。它需要两个参数:第一个是新列的名称字符串,第二个是新列本身。
新列必须是 Column 类的对象。创建这样的对象非常简单,只需通过 df.colName 从 DataFrame 中取出一列即可。
更新 Spark DataFrame 与在 pandas 中的做法有所不同,因为 Spark DataFrame 是不可变的(immutable)。也就是说,它不能被原地修改,因此列也无法就地更新。
因此,这些方法都会返回一个新的 DataFrame。若要覆盖原 DataFrame,您必须像下面这样将返回的 DataFrame 重新赋值给原变量:
df = df.withColumn("newCol", df.oldCol + 1)
以上代码会创建一个与 df 具有相同列、并新增 newCol 的 DataFrame,其中每个条目等于对应的 oldCol 条目再加 1。
要覆盖一列,只需将该列名作为第一个参数传入即可!
请记住,工作区中已经有一个名为 spark 的 SparkSession。
本练习是课程的一部分
PySpark 基础
练习说明
- 使用
spark.table()方法并传入参数"flights",创建一个包含.catalog中flights表数据的 DataFrame。将其保存为flights。 - 使用
flights.show()展示flights的表头。请检查输出:air_time列包含航班时长(单位为分钟)。 - 更新
flights,新增名为duration_hrs的列,其中包含每次航班以小时计的时长(您需要用每小时的分钟数去除以air_time)。
交互式实操练习
通过完成这段示例代码来试试这个练习。
# Create the DataFrame flights
flights = spark.table(____)
# Show the head
____.____()
# Add duration_hrs
flights = flights.withColumn(____)