ÎncepețiÎncepe gratuit

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

Vezi cursul

Instrucțiuni pentru exercițiu

  • Folosește metoda spark.table() cu argumentul "flights" pentru a crea un DataFrame care conține valorile din tabelul flights în .catalog. Salvează-l ca flights.
  • Afișează primele rânduri din flights folosind flights.show(). Verifică rezultatul: coloana air_time conține durata zborului în minute.
  • Actualizează flights pentru a include o coloană nouă numită duration_hrs, care să conțină durata fiecărui zbor în ore (va trebui să împarți air_time la 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(____)
Editează și rulează codul