Was ist die Spark DataFrame API?

Apache Spark ist eine verteilte Verarbeitungs-Engine, die strukturierte, semistrukturierte und unstrukturierte Daten in großem Umfang verarbeiten kann. Entwickler verwenden die Spark API, um mit der Engine zu interagieren und verteilte Datenvorgänge zu definieren.


Das Herzstück der modernen Spark-Entwicklung ist die Spark DataFrame API. Ein DataFrame organisiert und verarbeitet verteilte Daten in benannten Spalten, ähnlich wie eine Tabelle in einer relationalen Datenbank oder einer Tabellenkalkulation. Sie bietet eine strukturierte, deklarative Schnittstelle zur Analyse riesiger Datasets und optimiert gleichzeitig automatisch die physische Ausführung im Hintergrund.

Spark API: Architektur und Ausführung

Um zu verstehen, wie die Spark API Ihren Code in einem verteilten Netzwerk ausführt, ist es wichtig, die grundlegenden Komponenten der Laufzeitarchitektur zu kennen:


  • Treiber: Der zentrale Koordinator der Spark-Anwendung. Es liest Ihren Code, übersetzt deklarative Vorgänge in logische Ausführungspläne und plant Aufgaben auf den Worker-Knoten.
  • Cluster-Manager: Der Ressourcenzuweiser (z. B. Kubernetes oder eigenständige Ressourcenmanager), der CPU-, Arbeitsspeicher- und Netzwerkressourcen im gesamten Cluster bereitstellt.
  • Executors: Die Worker-Instanzen, die auf Clusterknoten ausgeführt werden. Sie erhalten Aufgabenanweisungen vom Treiber, führen Datenverarbeitungsvorgänge lokal im Arbeitsspeicher aus und geben Ergebnisse zurück oder schreiben sie in den Speicher.
  • DAGs, Phasen und Aufgaben: Spark führt Datentransformationen nicht zeilenweise aus. Stattdessen kompiliert der Treiber Ihren DataFrame-Code in einen gerichteten azyklischen Graphen (Directed Acyclic Graph, DAG) physischer Operatoren. Der Treiber gruppiert diese Operatoren in umfassende Ausführungsphasen, die als Phasen bezeichnet werden und durch Grenzen für die Datenumverteilung nach dem Zufallsprinzip unterteilt sind. Diese Phasen werden in einzelne Aufgaben unterteilt, die zur parallelen Ausführung an die Executors verteilt werden.

Was ist der Unterschied zwischen Spark-RDDs und DataFrames?

Bei der Entwicklung verteilter Datenpipelines können Entwickler zwischen zwei primären Programmierschnittstellen wählen:


  • Spark RDD (Resilient Distributed Dataset): Die ursprüngliche, Low-Level-Spark-API. Sie stellt eine unveränderliche, fehlertolerante Sammlung von JVM-Objekten dar, die über einen Cluster verteilt sind. Da RDDs beliebige Java-Objekte enthalten, kann die Ausführungs-Engine ihre interne Struktur nicht prüfen. Dies zwingt Entwickler dazu, Low-Level-Transformationslogik manuell zu schreiben und zu optimieren, was oft zu einem erheblichen Overhead bei der automatischen Speicherbereinigung und einer langsamen Serialisierung führt.
  • Spark DataFrame: Der moderne Standard für die Verarbeitung strukturierter Daten. DataFrames organisieren Daten in Zeilen und Spalten, die durch ein striktes Schema definiert sind. Da die Ausführungs-Engine die Datentypen und die Struktur des Datasets versteht, kann sie den Code vor der Ausführung automatisch optimieren.

Warum sollten Sie DataFrames anstelle von RDDs verwenden?

Für fast alle Data-Engineering- und Data-Science-Anwendungsfälle ist die DataFrame API die bevorzugte Wahl. RDDs sind für Szenarien reserviert, in denen Sie unstrukturierte Rohdaten (z. B. Binärdateien oder Medienstreams) mit benutzerdefinierten Java-Objekten bearbeiten müssen. DataFrames erfordern weniger Code, werden automatisch schneller ausgeführt und nutzen die Off-Heap-Binärspeicherverwaltung, um Engpässe bei der automatischen Speicherbereinigung der JVM zu umgehen.

So verwenden Datenteams die Spark API

Die Spark DataFrame API vereinfacht die rechenintensiven Aufgaben der Verarbeitung, Bereinigung und Vorbereitung großer Datenmengen.


  • Data Engineers: Engineers verwenden die DataFrame API, um robuste, fehlertolerante Datenpipelines zu erstellen. Sie extrahieren Rohdaten, wenden Schemas an, führen strukturierte Transformationen aus und schreiben konforme Tabellen in Data Lakes oder Cloud Data Warehouses. Die deklarative Syntax ermöglicht es ihnen, täglich Terabyte an Daten zu verarbeiten und gleichzeitig die Codekomplexität und den Wartungsaufwand zu minimieren.
  • Data Scientists: Data Scientists verwenden die DataFrame API, um riesige Datasets zu untersuchen und vorzubereiten, die die Arbeitsspeicherlimits einer einzelnen Maschine überschreiten. Durch die Verteilung von Datenpartitionen auf Worker-Knoten können sie explorative Datenanalysen durchführen, Nullwerte bereinigen und Features in großem Umfang entwickeln, was die Zeit bis zur Erkenntnisgewinnung verkürzt.

Kernbibliotheken, die auf der DataFrame API basieren

Die Spark DataFrame API dient als programmatische Grundlage für die erweiterten Verarbeitungsbibliotheken von Spark:


  • Spark SQL: Mit diesem Modul können Sie ANSI-SQL-Abfragen direkt für Spark-Datasets ausführen. Sie können DataFrames als temporäre Ansichten abfragen und so deklarativen Python- oder Scala-Code nahtlos mit Standard-SQL-Abfragen kombinieren.
  • Strukturiertes Streaming: Diese Engine verarbeitet kontinuierliche Streams von Echtzeitdaten. Es verwendet dieselben DataFrame-API-Befehle wie die statische Batchverarbeitung und übernimmt automatisch die Mikrostapelverarbeitung, Fehlertoleranz und Zustandsverwaltung.
  • MLlib: Die verteilte Machine-Learning-Bibliothek von Spark verwendet DataFrames, um Trainings-Datasets zu verwalten und vorzubereiten. Sie bietet integrierte, skalierbare Algorithmen für Klassifizierung, Regression, Clustering und kollaboratives Filtern.

Vorteile der Spark DataFrame API

Deklarative Optimierung

DataFrames verwenden einen integrierten Abfrageoptimierer namens Catalyst Optimizer. Wenn Sie DataFrame-Code schreiben, schreibt der Optimierer automatisch den physischen Ausführungsplan um und wendet Filter-Pushdown, Projektionsbereinigung und Join-Neuordnung an, um so schnell wie möglich ohne manuelle Abstimmung ausgeführt zu werden.

Sprachliche Flexibilität

Unabhängig davon, ob Ihr Team Code in Python (PySpark), Scala, Java oder R schreibt, sorgt die DataFrame API für eine identische Leistung. Der zugrunde liegende Ausführungsplan wird innerhalb derselben JVM-unabhängigen Ausführungsebene kompiliert und optimiert.

Effiziente Speicherverwaltung

DataFrames nutzt die Tungsten-Ausführungs-Engine, um Daten in einem hochkomprimierten, Off-Heap-Binärformat zu speichern. Dadurch wird der Overhead für die Erstellung von JVM-Objekten eliminiert und verhindert, dass Executor-Threads durch Pausen bei der automatischen Speicherbereinigung eingefroren werden.

Spark DataFrame API in Google Cloud verwenden

Die Ausführung von Open-Source-Spark erfordert herkömmlicherweise eine manuelle Clusterkonfiguration, die Verwaltung von Softwareversionen und komplexes Netzwerk-Peering. Google Cloud vereinfacht dies durch den Managed Service for Apache Spark, der die verteilte Ausführung in eine vollständig verwaltete, unternehmensreife Datenplattform verwandelt.


Datenteams führen Spark DataFrame-Arbeitslasten in Google Cloud mit den folgenden Funktionen aus:


  • Serverlose und verwaltete Clusterbereitstellung: Mit Google Cloud können Sie das Ausführungsmodell auswählen, das Ihren betrieblichen Anforderungen entspricht. Sie können DataFrame-Code im serverlosen Modus ausführen, um Batchjobs direkt zu senden. Dabei zahlen Sie nur für die genauen Sekunden der Laufzeit, da die Ressourcen automatisch skaliert werden. Alternativ können Sie hochgradig anpassbare, persistente verwaltete Cluster für kontinuierliche Arbeitslasten bereitstellen.
  • Optimierte BigQuery-Konnektivität: Der Spark-BigQuery-Connector umgeht zeilenorientierte JVM-Übergänge, indem er BigQuery-Daten direkt im nativen Apache Arrow-Format verarbeitet. Data Scientists können außerdem eine integrierte Notebook-Oberfläche verwenden, um PySpark-DataFrame-Code und SQL-Abfragen für dasselbe verwaltete Dataset auszuführen, ohne die Umgebung wechseln zu müssen.
  • Vektorisierte native C++-Ausführung: Wenn Sie Spark-Arbeitslasten in Google Cloud ausführen, können Sie die Lightning Engine für Ihre serverlosen Batches oder verwalteten Cluster aktivieren. Diese native C++-Abfrageausführungs-Engine kompiliert physische Abfragepläne mithilfe von Velox und Gluten direkt in native Anweisungen. Durch die Umgehung des JVM-Volcano-Iterator-Modells und der Engpässe bei der automatischen Speicherbereinigung werden Ihre DataFrame- und Spark SQL-Arbeitslasten ohne Codeänderungen um das bis zu 4,9-Fache beschleunigt.


Gleich loslegen

Profitieren Sie von einem Guthaben über 300 $, um Google Cloud und mehr als 20 „Immer kostenlos“ Produkte kennenzulernen.

Google Cloud