Crearea coloanelor
În acest capitol, vei învăța cum să folosești metodele definite de clasa DataFrame din Spark pentru a efectua operații comune asupra datelor.
Să vedem cum se realizează operații la nivel de coloană. În Spark, poți face acest lucru cu metoda .withColumn(), care primește două argumente. Primul este un șir de caractere cu numele noii coloane, iar al doilea este coloana în sine.
Noua coloană trebuie să fie un obiect de clasa Column. Crearea unui astfel de obiect este la fel de simplă ca extragerea unei coloane din DataFrame-ul tău folosind df.colName.
Actualizarea unui DataFrame Spark este ușor diferită față de lucrul cu pandas, deoarece DataFrame-ul Spark este immutable (imuabil). Asta înseamnă că nu poate fi modificat direct, deci coloanele nu pot fi actualizate în loc.
Prin urmare, toate aceste metode returnează un DataFrame nou. Pentru a suprascrie DataFrame-ul original, trebuie să reatribui DataFrame-ul returnat, astfel:
df = df.withColumn("newCol", df.oldCol + 1)
Codul de mai sus creează un DataFrame cu aceleași coloane ca df, plus o coloană nouă, newCol, în care fiecare valoare este egală cu valoarea corespunzătoare din oldCol, plus unu.
Pentru a suprascrie o coloană existentă, transmite pur și simplu numele acelei coloane ca prim argument!
Reține că un SparkSession numit spark este deja disponibil în spațiul tău de lucru.
Acest exercițiu face parte din cursul
Fundamente PySpark
Instrucțiuni pentru exercițiu
- Folosește metoda
spark.table()cu argumentul"flights"pentru a crea un DataFrame care conține valorile din tabelulflightsîn.catalog. Salvează-l caflights. - Afișează primele rânduri din
flightsfolosindflights.show(). Verifică rezultatul: coloanaair_timeconține durata zborului în minute. - Actualizează
flightspentru a include o coloană nouă numităduration_hrs, care să conțină durata fiecărui zbor în ore (va trebui să împarțiair_timela numărul de minute dintr-o oră).
Exercițiu interactiv practic
Încearcă acest exercițiu completând acest cod de exemplu.
# Create the DataFrame flights
flights = spark.table(____)
# Show the head
____.____()
# Add duration_hrs
flights = flights.withColumn(____)