Dataflow wurde entwickelt, um große Datenverarbeitungs-Pipelines auszuführen, indem die Arbeit auf einen verwalteten Pool von Compute-Instanzen verteilt wird. Wenn Sie wissen, wie Dataflow die Verarbeitung parallelisiert, können Sie effiziente Pipelines entwerfen, Leistungsengpässe vermeiden und Ressourcenkosten optimieren.
Auf dieser Seite wird erläutert, wie Dataflow die Datenverarbeitung parallelisiert, wie die Ausführung verwaltet und skaliert wird, welche Faktoren die Parallelisierung einschränken und welche Techniken Sie verwenden können, um den Durchsatz von Pipelines zu optimieren.
Parallelitätsmodelle: horizontal und vertikal
Dataflow erreicht Parallelität mit zwei sich ergänzenden Strategien:
Horizontale Parallelität:Pipelinedaten werden partitioniert und gleichzeitig auf mehreren Worker-Instanzen (virtuellen Maschinen) verarbeitet. Mit horizontalem Autoscaling kann Dataflow die Größe des Worker-Pools automatisch an die Anforderungen der Arbeitslast anpassen. Standardmäßig legt Dataflow ein Ressourcenlimit von 4.000 Workern pro Job fest, das mit Kontingentanfragen angepasst werden kann.
Vertikale Parallelität: Mehrere CPU-Kerne und Threads in einer einzelnen Worker-Instanz verarbeiten Pipelinedaten gleichzeitig. Auf jeder Worker-VM werden Worker-Prozesse und Harness-Threads ausgeführt, um die verfügbaren Rechenressourcen zu nutzen. Mit Dynamic Thread Scaling kann Dataflow die Anzahl der aktiven Threads pro Worker in Batchpipelines basierend auf der CPU-Auslastung und dem verfügbaren Arbeitsspeicher anpassen. In Dataflow Prime skaliert vertikales Autoscaling den für Worker zugewiesenen Arbeitsspeicher und die Rechenleistung dynamisch.
Arbeitseinheiten und Ausführungshierarchie
Um die Verarbeitung auf Worker und Threads zu verteilen, unterteilt Dataflow Apache Beam-Pipelines in separate Arbeitseinheiten:
- PCollections und Partitionen:Ein
PCollectionstellt einen verteilten Datensatz dar. Bei begrenzten Daten (Batchpipelines) teilt Dataflow das Dataset in Splits oder Shards auf. Bei unbegrenzten Daten (Streaming-Pipelines) werden Daten kontinuierlich empfangen und als Nachrichten oder Streampartitionen aufgenommen. - Gruppierungen:Dataflow gruppiert Elemente in beliebigen Gruppierungen für die Verarbeitung durch eine
DoFn. Eine Gruppierung ist die Einheit für Fehler und Wiederholungen. Wenn bei der Verarbeitung eines Elements eine unbehandelte Ausnahme ausgelöst wird, wird die gesamte Gruppierung noch einmal verarbeitet. Vorgänge mit hohem Arbeitsspeicherverbrauch können den Arbeitsspeicher der Worker belasten und zu Arbeitsspeicherfehlern führen. - Phasen und Zusammenführen von Schritten:Bei der Diagrammoptimierung fasst Dataflow benachbarte Transformationen in zusammengeführte Ausführungsphasen zusammen, um den Overhead der Zwischenmaterialisierung von Daten zu vermeiden. In einer zusammengeführten Phase werden Elemente in einer engen Ausführungsschleife in einem einzelnen Thread verarbeitet, bevor sie an die nächste Phase oder Shuffle-Grenze übergeben werden.
Weitere Informationen zur Pipelineübersetzung und zur Grafikerstellung finden Sie unter Pipeline-Lebenszyklus.
Verwaltete Parallelität und Autoscaling
Standardmäßig verwaltet Dataflow die Pipelineparallelität automatisch, ohne dass eine manuelle Partitionsoptimierung erforderlich ist. Dies geschieht auf folgende Weise:
- Horizontales Autoscaling:
- Batch-Pipelines:Die Gesamtzahl der geschätzten verbleibenden Aufgaben, der Quell-Backlog und die CPU-Auslastung werden ausgewertet, um den Worker-Pool zu skalieren und den Job schnell und kostengünstig abzuschließen.
- Streamingpipelines:Analysiert die Systemlatenz, die Rückstandsgröße und die CPU-Auslastung, um Worker bei Durchsatzspitzen hochzuskalieren und bei geringem Traffic herunterzuskalieren. Weitere Informationen finden Sie unter Horizontales Streaming-Autoscaling optimieren.
- Dynamischer Arbeitsausgleich (Dynamic Work Rebalancing, DWR): In Batchpipelines überwacht Dataflow den Fortschritt einzelner Worker-Aufgaben. Wenn ein Worker frühzeitig fertig ist oder ein anderer Worker aufgrund von Datenverzerrung (Nachzügler) hinterherhinkt, teilt Dataflow die verbleibende, nicht verarbeitete Arbeit des langsamen Workers dynamisch auf und weist sie einem inaktiven Worker neu zu. Weitere Informationen finden Sie unter Dynamischer Work-Ausgleich.
- Dynamische Thread-Skalierung:Bei Batchpipelines, die den Portable Runner verwenden, wird die Anzahl der Threads für die gleichzeitige Verarbeitung pro Worker automatisch an die CPU-Auslastung und den verfügbaren Arbeitsspeicher angepasst. Weitere Informationen finden Sie unter Dynamische Thread-Skalierung.
- Vertikales Autoscaling:In Dataflow Prime skaliert Dataflow den Arbeitsspeicher und die Rechenressourcen von Workern dynamisch, um Speicherplatzmangel zu vermeiden und die Ressourcennutzung zu optimieren. Weitere Informationen finden Sie unter Vertikales Autoscaling.
Faktoren, die die Parallelität einschränken
Eine Pipeline erreicht möglicherweise nicht die erwartete Parallelität, wenn die folgenden Datenmerkmale oder das Pipelinegrafikdesign vorliegen:
Nicht aufteilbare Eingabequellen
Wenn eine Eingabequelle nicht in unabhängige Bereiche aufgeteilt werden kann, muss Dataflow die Quelle sequenziell mit einem einzelnen Arbeitsthread lesen:
- Nicht aufteilbare Dateikomprimierung:Formate wie
.gz(gzip) oder.bzip2(ohne Indexierung) können nicht parallel von beliebigen Byte-Offsets gelesen werden. Das Lesen einer einzelnen großen komprimierten Datei beschränkt die Erfassungsphase auf einen einzelnen Thread, bis die Daten dekomprimiert und neu verteilt werden. - Lösung:Speichern Sie Daten in aufteilbaren Dateiformaten wie Parquet, Avro oder Snappy-komprimierten Formaten oder teilen Sie Eingabedaten in mehrere kleinere Dateien in Cloud Storage auf.
Schrittzusammenführung und hohes Fan-Out
Durch die Zusammenführung von Schritten wird die Leistung verbessert, da der Serialisierungsaufwand verringert wird. Sie kann jedoch unbeabsichtigt die Parallelität einschränken und die Arbeitsspeicherlast erhöhen, wenn ein Schritt mit geringer Parallelität eine große Anzahl von Ausgabeelementen erzeugt (ein Vorgang mit hohem Fan-out):
- Beispiel:Eine Quelle liest fünf Dateien und wird mit einer
FlatMap-Transformation zusammengeführt, die 1.000.000 Ausgabeelemente erzeugt. Wenn dieFlatMap-Transformation mit nachgelagerten Transformationen zusammengeführt wird, werden alle 1.000.000 Elemente weiterhin auf maximal fünf Worker-Threads ausgeführt, was den nachgelagerten Durchsatz erheblich einschränkt. Wenn Zwischen-Transformationen vor dem Committing im Arbeitsspeicher erheblich erweitert werden, können große Bundles den verfügbaren Worker-Arbeitsspeicher erschöpfen. - Lösung:Fügen Sie eine
Redistribute-Transformation (oder eine klassischeReshuffle-Transformation) zwischen dem Schritt mit hohem Fan-Out und den nachgelagerten Transformationen ein, um die Zusammenführung aufzuheben und die Arbeit auf den Worker-Pool zu verteilen. Informationen zum Debugging von speicherbezogenen Problemen finden Sie unter Fehlerbehebung bei Out of Memory-Fehlern.
Tastaturbelegung und Hotkeys
Bei Aggregationsvorgängen (GroupByKey, CoGroupByKey, Combine.PerKey) werden Elemente nach dem zugehörigen Schlüssel gruppiert.
- Engpass durch übergroßen Schlüssel:Dataflow leitet alle Elemente mit demselben Schlüssel zur Aggregation an einen einzelnen Arbeitsthread weiter. Wenn ein einzelner Schlüssel einen großen Prozentsatz des gesamten Datasets enthält, wird dieser Worker zum Nachzügler und Upstream-Worker können einen Rückstau erfahren. Das kann beispielsweise ein Standardschlüssel für
nulloder ein Schlüssel für eine sehr beliebte Kategorie sein. - Lösung:
- Verwenden Sie nach Möglichkeit Combiners (
CombineFnoderCombine.PerKey) anstelle vonGroupByKey, damit Dataflow vor dem Shuffle partielle lokale Kombinationen ausführen kann. - Fügen Sie Hotkeys ein zufälliges Ganzzahlpräfix oder ‑suffix hinzu (Key Salting), um den Schlüsselbereich auf die Worker zu verteilen. Führen Sie dann eine Aggregation in der zweiten Phase durch, um die gesalzenen Ergebnisse zusammenzuführen.
- Verwenden Sie nach Möglichkeit Combiners (
Drosselung nachgelagerter Senken
Wenn Sie die Pipelineausgabe in externe Dienste wie Datenbanken oder Drittanbieter-APIs schreiben, kann eine hohe Parallelität das Zielsystem überlasten:
- Drosselung:Hunderte von Worker-Threads, die gleichzeitig Schreibaufrufe ausgeben, können zu Ratenbegrenzungsfehlern, Verbindungszeitüberschreitungen oder einer Beeinträchtigung der Datenbank führen.
- Lösung:
- Sie können die Parallelität von Schreibvorgängen einschränken, indem Sie Elemente mit
GroupByKeygruppieren oder Batching-Senken mit kontrollierter Parallelität verwenden. - Implementieren Sie clientseitigen exponentiellen Backoff und Wiederholungslogik in Senkenimplementierungen für
DoFn.
- Sie können die Parallelität von Schreibvorgängen einschränken, indem Sie Elemente mit
Optimierungsstrategien
Um die Parallelität in Ihren Dataflow-Jobs zu optimieren, sollten Sie die folgenden Ansätze in Betracht ziehen:
- Unerwünschte Zusammenführungen mit
Redistributeverhindern:Redistribute.arbitrarily(): Unterbricht die Zusammenführung von Schritten und verteilt Elemente gleichmäßig auf alle verfügbaren Worker.Redistribute.byKey(): Gleicht Schlüssel/Wert-Paare über Worker-Threads hinweg aus, wobei die Schlüssel-Lokalität beibehalten wird.- Implementierungsbeispiele finden Sie unter Zusammenführung verhindern.
- Nachzügler und Engpässe beobachten:Verwenden Sie die Ausführungsdetails der Google Cloud Console, um Phasen mit einer hohen Anzahl von Nachzüglern oder einem ins Stocken geratenen Fortschritt zu ermitteln:
Nächste Schritte
- Pipeline-Lebenszyklus
- Horizontales Autoscaling
- Informationen zum dynamischen Work-Ausgleich
- Best Practices für Dataflow-Pipelines
- Informationen zur Fehlerbehebung bei Out of Memory-Fehlern