Создание столбцов
В этой главе вы узнаете, как использовать методы класса DataFrame в Spark для выполнения типичных операций с данными.
Рассмотрим операции со столбцами. В Spark для этого используется метод .withColumn(), который принимает два аргумента: первый — строка с именем нового столбца, второй — сам новый столбец.
Новый столбец должен быть объектом класса Column. Создать его так же просто, как извлечь столбец из DataFrame с помощью df.colName.
Обновление Spark DataFrame несколько отличается от работы с pandas, поскольку Spark DataFrame является неизменяемым. Это означает, что его нельзя изменить напрямую, а значит, столбцы нельзя обновить на месте.
Поэтому все эти методы возвращают новый DataFrame. Чтобы перезаписать исходный DataFrame, необходимо переприсвоить результат следующим образом:
df = df.withColumn("newCol", df.oldCol + 1)
Приведённый код создаёт DataFrame с теми же столбцами, что и df, плюс новый столбец newCol, где каждое значение равно соответствующему значению из oldCol, увеличенному на единицу.
Чтобы перезаписать существующий столбец, просто передайте его имя в качестве первого аргумента!
Напомним: объект SparkSession с именем spark уже доступен в вашем рабочем пространстве.
Это упражнение является частью курса
Основы PySpark
Инструкции к упражнению
- Используйте метод
spark.table()с аргументом"flights", чтобы создать DataFrame на основе таблицыflightsиз.catalog. Сохраните результат в переменнуюflights. - Выведите первые строки
flightsс помощьюflights.show(). Обратите внимание на вывод: столбецair_timeсодержит продолжительность рейса в минутах. - Обновите
flights, добавив новый столбецduration_hrsс продолжительностью каждого рейса в часах (для этого разделитеair_timeна количество минут в одном часе).
Интерактивное практическое упражнение
Попробуйте выполнить это упражнение, дополнив этот пример кода.
# Create the DataFrame flights
flights = spark.table(____)
# Show the head
____.____()
# Add duration_hrs
flights = flights.withColumn(____)