Diese Seite bietet eine Übersicht über den Pipelinelebenszyklus – vom Pipelinecode zu einem Dataflow-Job.
Auf dieser Seite werden die folgenden Konzepte erläutert:
- Was eine Ausführungsgrafik ist und wie eine Apache Beam-Pipeline zu einem Dataflow-Job wird.
- So verarbeitet Dataflow Fehler.
- So parallelisiert und verteilt Dataflow die Verarbeitungslogik der Pipeline automatisch an die Worker, die Ihren Job ausführen.
- Joboptimierungen, die Dataflow vornehmen könnte
Ausführungsgrafik
Wenn Sie das Dataflow-Pipeline ausführen, erstellt Dataflow anhand des für das Pipeline-Objekt verwendeten Codes eine Ausführungsgrafik. Diese enthält alle Transformationen sowie die zugehörigen Verarbeitungsfunktionen, wie etwa DoFn-Objekte. Dies ist die Ausführungsgrafik der Pipeline. Die Phase wird als Grafikerstellungszeit bezeichnet.
Während die Grafik erstellt wird, führt Apache Beam den Code lokal vom Haupteinstiegspunkt des Pipelinecodes aus, stoppt bei den Aufrufen einer Quelle, einer Senke oder eines Transformationsschritts und wandelt diese Aufrufe in Knoten der Grafik um.
Folglich wird ein Stück Code, das sich im Einstiegspunkt einer Pipeline (Java- und Go-Methode main oder die oberste Ebene eines Python-Skripts) befindet, lokal auf dem Computer ausgeführt, auf dem die Pipeline ausgeführt wird. Derselbe Code, der in einer Methode eines DoFn-Objekts deklariert ist, wird in den Dataflow-Workern ausgeführt.
Das in den Apache Beam SDKs enthaltene WordCount enthält eine Reihe von Transformationen zum Lesen, Extrahieren, Zählen, Formatieren und Schreiben der einzelnen Wörter in einer Textsammlung. Außerdem wird zu jedem Wort angegeben, wie oft es vorkommt. Das folgende Diagramm zeigt, wie die Transformationen in der Pipeline von WordCount zu einer Ausführungsgrafik erweitert werden:

Abbildung 1: Ausführungsgrafik des Beispiels WordCount
Die Ausführungsgrafik weicht oft von der Reihenfolge ab, in der Sie die Transformationen beim Erstellen der Pipeline angegeben haben. Der Unterschied ist, dass der Dataflow-Dienst verschiedene Optimierungen und Zusammenlegungen in der Ausführungsgrafik vornimmt, bevor diese in verwalteten Cloudressourcen ausgeführt wird. Der Dataflow-Dienst berücksichtigt beim Ausführen der Pipeline Datenabhängigkeiten. Die Schritte ohne Datenabhängigkeiten zwischen ihnen können jedoch in beliebiger Reihenfolge ausgeführt werden.
Wenn Sie die von Dataflow für die Pipeline erstellte, nicht optimierte Ausführungsgrafik sehen möchten, wählen Sie den Job auf der Dataflow-Monitoring-Oberfläche aus. Weitere Informationen zum Aufrufen von Jobs finden Sie unter Dataflow-Monitoring-Oberfläche verwenden.
Apache Beam prüft während der Grafikerstellung, ob alle Ressourcen, auf die die Pipeline verweist (wie Cloud Storage-Buckets, BigQuery-Tabellen und Pub/Sub-Themen und -Abos), tatsächlich vorhanden und zugänglich sind. Die Prüfung erfolgt über Standard-API-Aufrufe der jeweiligen Dienste. Daher ist es wichtig, dass das zum Ausführen einer Pipeline verwendete Nutzerkonto eine funktionierende Verbindung zu den erforderlichen Diensten und die Berechtigung zum Aufrufen von deren APIs hat. Bevor die Pipeline an den Dataflow-Dienst gesendet wird, prüft Apache Beam auch, ob weitere Fehler vorliegen, und sorgt dafür, dass die Pipelinegrafik keine ungültigen Vorgänge enthält.
Die Ausführungsgrafik wird dann in das JSON-Format übersetzt und an den Dataflow-Dienstendpunkt übertragen.
Der Dataflow-Dienst validiert dann die JSON-Ausführungsgrafik. Durch die Validierung wird die Grafik zu einem Job im Dataflow-Dienst. Sie können den Job, die Ausführungsgrafik, den Status und die Loginformationen auf der Dataflow-Monitoring-Oberfläche aufrufen.
Java
Der Dataflow-Dienst sendet eine Antwort an den Computer, auf dem Sie das Dataflow-Programm ausführen. Diese Antwort wird im Objekt DataflowPipelineJob gekapselt, das den jobId Ihres Dataflow-Jobs enthält.
Verwenden Sie jobId für Monitoring, Verfolgung und Fehlerbehebung Ihres Jobs mithilfe der Dataflow-Monitoring-Oberfläche und der Dataflow-Befehlszeile.
Weitere Informationen finden Sie in der API-Referenz zu "DataflowPipelineJob".
Python
Der Dataflow-Dienst sendet eine Antwort an den Computer, auf dem Sie das Dataflow-Programm ausführen. Diese Antwort wird im Objekt DataflowPipelineResult gekapselt, das den job_id Ihres Dataflow-Jobs enthält.
Verwenden Sie job_id für Monitoring, Verfolgung und Fehlerbehebung Ihres Jobs mithilfe der Dataflow-Monitoring-Oberfläche und der Dataflow-Befehlszeile.
Go
Der Dataflow-Dienst sendet eine Antwort an den Computer, auf dem Sie das Dataflow-Programm ausführen. Diese Antwort wird im Objekt dataflowPipelineResult gekapselt, das den jobID Ihres Dataflow-Jobs enthält.
Verwenden Sie jobID für Monitoring, Verfolgung und Fehlerbehebung Ihres Jobs mithilfe der Dataflow-Monitoring-Oberfläche und der Dataflow-Befehlszeile.
Die Grafik wird auch erstellt, wenn Sie die Pipeline lokal ausführen. Es erfolgt jedoch keine Übersetzung in das JSON-Format und keine Übertragung an den Dienst. Die Grafik wird stattdessen lokal auf demselben Computer ausgeführt, auf dem Sie Dataflow gestartet haben. Weitere Informationen finden Sie unter PipelineOptions für die lokale Ausführung konfigurieren.
Verarbeitung von Fehlern und Ausnahmen
Die Pipeline kann während der Datenverarbeitung Ausnahmen ausgeben. Einige dieser Fehler sind temporär, wie etwa vorübergehende Probleme beim Zugriff auf einen externen Dienst. Andere Fehler sind permanent, z. B. Fehler, die durch beschädigte oder nicht parsingfähige Eingabedaten verursacht werden, oder Nullzeiger während der Berechnung.
Dataflow verarbeitet Elemente in beliebigen Gruppierungen. Sollte für eines der Elemente in der Gruppierung ein Fehler ausgegeben werden, wird die gesamte Gruppierung noch einmal verarbeitet. Im Batchmodus wird die Verarbeitung von Gruppierungen mit einem fehlerhaften Element viermal wiederholt. Wenn eine Gruppierung viermal fehlgeschlagen ist, fällt die gesamte Pipeline aus. Im Streamingmodus wird die Verarbeitung einer Gruppierung mit einem fehlerhaften Element unendlich oft wiederholt. Dies kann zur permanenten Blockierung der Pipeline führen.
Bei der Verarbeitung im Batchmodus tritt möglicherweise erst eine größere Anzahl einzelner Fehler auf, bevor ein kompletter Pipelinejob fehlschlägt. Dies ist der Fall, wenn eine gegebene Gruppierung nach vier Wiederholversuchen fehlschlägt. Wenn die Pipeline beispielsweise versucht, 100 Gruppierungen zu verarbeiten, kann Dataflow mehrere Hundert einzelne Fehler generieren, bevor eine Gruppierung viermal ausfällt und die Verarbeitung beendet wird.
Fehler bei Start-Workern, z. B. die Nichtinstallation von Paketen auf den Workern, sind temporär. Dieses Szenario führt zu unbeschränkten Wiederholungen und kann zur permanenten Blockierung der Pipeline führen.
Parallelisierung und Verteilung
Der Dataflow-Dienst parallelisiert und verteilt die Verarbeitungslogik in Ihrer Pipeline automatisch auf Worker und Threads.
Dataflow verwendet die Abstraktionen im Programmiermodell, um Funktionen zur parallelen Verarbeitung darzustellen. Beispiel: ParDo-Transformationen führen dazu, dass Dataflow den Verarbeitungscode, der durch DoFn-Objekte dargestellt wird, auf mehrere Worker verteilt, die gleichzeitig ausgeführt werden sollen.
Dataflow unterstützt zwei sich ergänzende Dimensionen der Parallelität:
- Horizontale Parallelität:Pipelinedaten werden auf mehrere Worker-Instanzen aufgeteilt und gleichzeitig verarbeitet. Die Verwaltung erfolgt dynamisch über horizontales Autoscaling.
- Vertikale Parallelität:Pipelinedaten werden auf mehreren CPU-Kernen und Threads auf jeder Worker-VM verarbeitet. Die Verwaltung erfolgt über dynamische Thread-Skalierung und vertikales Autoscaling.
Dataflow verwaltet automatisch die Jobparallelität, sorgt für einen dynamischen Arbeitsausgleich und optimiert den Ausführungsgraphen durch Fusion und Combine-Optimierungen. Faktoren wie nicht aufteilbare Datenquellen, Schritte mit hohem Fan-Out, Schlüsselabweichung und Downstream-Senkenlimits können die Parallelität von Pipelines einschränken.
Eine ausführliche Anleitung dazu, wie Dataflow Daten partitioniert, Worker skaliert und Engpässe bei der Parallelität behebt, finden Sie unter Parallelität in Dataflow.
Zusammenführung optimieren
Nachdem die JSON-Form der Ausführungsgrafik der Pipeline validiert wurde, ändert der Dataflow-Dienst die Grafik möglicherweise, um Optimierungen vorzunehmen.
Eine Optimierung kann das Zusammenführen mehrerer Schritte oder Transformationen in der Pipeline-Ausführungsgrafik zu einem Schritt beinhalten. Durch das Zusammenführen von Schritten muss der Dataflow-Dienst nicht jede dazwischenliegende PCollection in der Pipeline erfassen, was im Hinblick auf den Speicher- und Verarbeitungsaufwand kostspielig werden kann.
Obwohl alle beim Erstellen der Pipeline angegebenen Transformationen für den Dienst ausgeführt werden, sollten die Transformationen möglicherweise in einer anderen Reihenfolge oder als Teil einer größeren Transformation ausgeführt werden, um eine möglichst effektive Ausführung Ihrer Pipeline zu gewährleisten. zusammengeführte Transformation. Dabei berücksichtigt der Dataflow-Dienst Datenabhängigkeiten zwischen den Schritten in der Ausführungsgrafik. Ansonsten können die Schritte jedoch in beliebiger Reihenfolge ausgeführt werden.
Beispiel einer Zusammenführung
Das folgende Diagramm zeigt, wie der Dataflow-Dienst die Ausführungsgrafik des WordCount-Beispiels aus dem Apache Beam SDK for Java optimieren und zusammenführen kann, um die Ausführung effizienter zu gestalten:

Abbildung 2: Für das WordCount-Beispiel optimierte Ausführungsgrafik
Zusammenführung vermeiden
In einigen Fällen kann Dataflow eine falsche optimale Zusammenführung von Vorgängen in der Pipeline erraten. Dies kann die Fähigkeit von Dataflow einschränken, alle verfügbaren Worker zu verwenden. In solchen Fällen können Sie Dataflow einen Hinweis geben, die Daten neu zu verteilen, indem Sie eine Redistribute-Transformation verwenden.
Rufen Sie eine der folgenden Methoden auf, um eine Redistribute-Transformation hinzuzufügen:
Redistribute.arbitrarily: Gibt an, dass die Daten wahrscheinlich unausgewogen sind. Dataflow wählt den besten Algorithmus zum Umverteilen der Daten aus.Redistribute.byKey: Gibt an, dass einPCollectionvon Schlüssel/Wert-Paaren wahrscheinlich unausgewogen ist und basierend auf den Schlüsseln neu verteilt werden sollte. Normalerweise werden alle Elemente eines einzelnen Schlüssels im selben Arbeitsthread platziert. Die gemeinsame Platzierung von Schlüsseln wird jedoch nicht garantiert und die Elemente werden unabhängig voneinander verarbeitet.
Wenn Ihre Pipeline eine Redistribute-Transformation enthält, verhindert Dataflow in der Regel das Zusammenführen der Schritte vor und nach der Redistribute-Transformation und mischt die Daten, sodass die Schritte nach der Redistribute-Transformation eine optimale Parallelität aufweisen.
Zusammenführung überwachen
Sie können über die gcloud CLI oder die API in der Google Cloud -Console auf die optimierten Graphen und die zusammengeführten Phasen zugreifen.
Console
Öffnen Sie auf dem Tab Ausführungsdetails des Dataflow-Jobs den Phasenworkflow in der Diagrammansicht, um die zusammengeführten Phasen und Schritte in der Console aufzurufen.
Klicken Sie im Diagramm auf die zusammengeführte Phase, um die Komponentenschritte zu sehen, die für eine Phase zusammengeführt wurden. Im Bereich Phaseninformationen werden in der Zeile Komponentenschritte die zusammengeführten Phasen angezeigt. Manchmal werden Teile einer einzelnen zusammengesetzten Transformation in mehrere Phasen zusammengeführt.
gcloud
Wenn Sie mit der gcloud CLI auf die optimierte Grafik und die zusammengeführten Phasen zugreifen möchten, führen Sie den folgenden gcloud-Befehl aus:
gcloud dataflow jobs describe --full JOB_ID --format json
Ersetzen Sie JOB_ID durch die ID Ihres Dataflow-Jobs.
Leiten Sie die Ausgabe des Befehls gcloud an jq weiter, um die relevanten Bits zu extrahieren:
gcloud dataflow jobs describe --full JOB_ID --format json | jq '.pipelineDescription.executionPipelineStage\[\] | {"stage_id": .id, "stage_name": .name, "fused_steps": .componentTransform }'
Um die Beschreibung der zusammengeführten Phasen in der Ausgabeantwortdatei anzusehen, rufen Sie im ComponentTransform-Array dasExecutionStageSummary-Objekt auf.
API
Wenn Sie über die API auf die optimierte Grafik und die zusammengeführten Phasen zugreifen möchten, rufen Sie project.locations.jobs.get auf.
Um die Beschreibung der zusammengeführten Phasen in der Ausgabeantwortdatei anzusehen, rufen Sie im ComponentTransform-Array dasExecutionStageSummary-Objekt auf.
Optimierung kombinieren
Zusammenführungen stellen bei umfangreichen Datenverarbeitungen ein wichtiges Konzept dar.
Dabei werden konzeptionell völlig unterschiedliche Daten vereint. Dies ist für Korrelationen extrem nützlich. Das Programmiermodell von Dataflow stellt die Aggregationsvorgänge in Form der Transformationen GroupByKey, CoGroupByKey und Combine dar.
Bei den Zusammenführungen durch Dataflow werden Daten des gesamten Datasets kombiniert, einschließlich Daten, die auf mehrere Worker verteilt sein können. Während der Zusammenführung ist es häufig am effizientesten, vor einer instanzübergreifenden Kombination so viele Daten wie möglich lokal zusammenzuführen. Wenn Sie eine Transformation vom Typ GroupByKey oder eine andere Aggregationstransformation anwenden, werden die Daten vom Dataflow-Dienst automatisch teilweise lokal kombiniert, bevor die Hauptgruppierung erfolgt.
Bei der Durchführung einer teilweisen oder mehrstufigen Kombination trifft der Dataflow-Dienst abhängig davon, ob die Pipeline mit Batch- oder Streamingdaten arbeitet, unterschiedliche Entscheidungen. Bei eingeschränkten Daten priorisiert der Dienst eine effiziente Vorgehensweise und kombiniert möglichst viel lokal. Bei uneingeschränkten Daten bevorzugt der Dienst eine geringere Latenz und führt möglicherweise keine teilweise Kombination durch, da dies die Latenz erhöhen kann.