Dataflow est conçu pour exécuter de grands pipelines de traitement de données en distribuant le travail sur un pool géré d'instances de calcul. Comprendre comment Dataflow parallélise le traitement vous aide à concevoir des pipelines efficaces, à éviter les goulots d'étranglement des performances et à optimiser les coûts des ressources.
Cette page explique comment Dataflow parallélise le traitement des données, comment il gère et met à l'échelle l'exécution, les facteurs courants qui limitent le parallélisme et les techniques que vous pouvez utiliser pour optimiser le débit du pipeline.
Modèles de parallélisme : horizontal et vertical
Dataflow atteint le parallélisme en utilisant deux stratégies complémentaires :
Parallélisme horizontal : les données du pipeline sont partitionnées et traitées simultanément sur plusieurs instances de nœud de calcul (machines virtuelles). Dataflow peut ajuster automatiquement la taille du pool de nœuds de calcul en fonction de la demande de charge de travail grâce à l'autoscaling horizontal. Par défaut, Dataflow définit une limite de ressources de 4 000 nœuds de calcul par job, qui peut être ajustée à l'aide des demandes de quota.
Parallélisme vertical : plusieurs cœurs et threads de processeur au sein d'une même instance de nœud de calcul traitent les données du pipeline simultanément. Chaque VM de nœud de calcul exécute des processus de nœud de calcul et exploite des threads de harnais pour utiliser les ressources de calcul disponibles. Avec le scaling dynamique de threads, Dataflow peut ajuster le nombre de threads actifs par nœud de calcul dans les pipelines par lot en fonction de l'utilisation du processeur et de la marge de mémoire. Dans Dataflow Prime, l'autoscaling vertical adapte de manière dynamique la mémoire et le calcul alloués aux nœuds de calcul.
Unités de travail et hiérarchie d'exécution
Pour distribuer le traitement entre les nœuds de calcul et les threads, Dataflow divise les pipelines Apache Beam en unités de travail distinctes :
- PCollections et partitions : une
PCollectionreprésente un ensemble de données distribué. Pour les données limitées (pipelines par lot), Dataflow divise l'ensemble de données en fractions ou en partitions. Pour les données illimitées (pipelines de traitement en flux continu), les données arrivent en continu et sont ingérées sous forme de messages ou de partitions de flux. - Ensembles : Dataflow regroupe les éléments dans des ensembles arbitraires pour qu'ils soient traités par un
DoFn. Un bundle est l'unité d'échec et de nouvelle tentative : si le traitement d'un élément génère une exception non gérée, l'ensemble du bundle est relancé. Les opérations qui consomment beaucoup de mémoire peuvent augmenter la pression sur la mémoire des nœuds de calcul et entraîner des erreurs de mémoire insuffisante. - Étapes et fusion d'étapes : lors de l'optimisation du graphique, Dataflow combine les transformations adjacentes en étapes d'exécution fusionnées pour éliminer la surcharge de la matérialisation des données intermédiaires. Dans une étape fusionnée, les éléments sont traités dans une boucle d'exécution étroite sur un seul thread avant d'être transmis à l'étape suivante ou à la limite de shuffle.
Pour en savoir plus sur la traduction de pipeline et la génération de graphiques, consultez Cycle de vie du pipeline.
Parallélisme géré et autoscaling
Par défaut, Dataflow gère automatiquement le parallélisme du pipeline sans nécessiter de réglage manuel des partitions, de la manière suivante :
- Autoscaling horizontal :
- Pipelines par lot : évalue le total du travail restant estimé, le backlog de la source et l'utilisation du processeur pour ajuster la taille du pool de nœuds de calcul afin de terminer le job rapidement et de manière économique.
- Pipelines de traitement en flux continu : ils analysent la latence du système, la taille du backlog et l'utilisation du processeur pour augmenter le nombre de nœuds de calcul lors des pics de débit et le réduire lors des périodes de faible trafic. Pour en savoir plus, consultez Régler l'autoscaling horizontal des flux.
- Rééquilibrage dynamique du travail (DWR) : dans les pipelines par lot, Dataflow surveille la progression des tâches de chaque nœud de calcul. Si un nœud de calcul termine son travail plus tôt que prévu ou si un autre nœud de calcul est en retard en raison d'un déséquilibre des données (nœuds de calcul lents), Dataflow divise dynamiquement le travail résiduel non traité du nœud de calcul lent et le réattribue à un nœud de calcul inactif. Pour en savoir plus, consultez Rééquilibrage dynamique du travail.
- Scaling dynamique de threads : dans les pipelines par lot qui utilisent le Portable Runner, ajuste automatiquement le nombre de threads de traitement simultanés par nœud de calcul en fonction de l'utilisation du processeur et de la marge de mémoire. Pour en savoir plus, consultez la section Scaling dynamique de threads.
- Autoscaling vertical : dans Dataflow Prime, Dataflow adapte de manière dynamique la mémoire et les ressources de calcul des nœuds de calcul pour éviter les erreurs de mémoire insuffisante et optimiser l'utilisation des ressources. Pour en savoir plus, consultez Autoscaling vertical.
Facteurs limitant le parallélisme
Un pipeline peut ne pas atteindre le parallélisme attendu en raison des caractéristiques de données ou de la conception du graphique de pipeline suivantes :
Sources d'entrée non fractionnables
Si une source d'entrée ne peut pas être divisée en plages indépendantes, Dataflow est obligé de lire la source de manière séquentielle avec un seul thread de nœud de calcul :
- Compression de fichiers non fractionnables : les formats tels que
.gz(gzip) ou.bzip2(sans indexation) ne peuvent pas être lus en parallèle à partir de décalages d'octets arbitraires. La lecture d'un seul fichier compressé volumineux limite l'étape d'ingestion à un seul thread jusqu'à ce que les données soient décompressées et redistribuées. - Résolution : stockez les données dans des formats de fichiers fractionnables (tels que Parquet, Avro ou les formats compressés avec Snappy) ou divisez les données d'entrée en plusieurs fichiers plus petits dans Cloud Storage.
Fusion d'étapes et distribution ramifiée élevée
La fusion des étapes améliore les performances en réduisant la surcharge de sérialisation, mais elle peut involontairement limiter le parallélisme et augmenter la pression sur la mémoire lorsqu'une étape à faible parallélisme produit un grand nombre d'éléments de sortie (une opération à "distribution ramifiée") :
- Exemple : Une source lit cinq fichiers et est fusionnée avec une transformation
FlatMapqui produit 1 000 000 d'éléments de sortie. Si la transformationFlatMapest fusionnée avec des transformations en aval, les 1 000 000 d'éléments continuent de s'exécuter sur un maximum de cinq threads de calcul, ce qui limite considérablement le débit en aval. De plus, si les transformations intermédiaires se développent considérablement en mémoire avant d'être validées, les grands bundles peuvent épuiser la mémoire disponible des nœuds de calcul. - Résolution : insérez une transformation
Redistribute(ouReshuffleclassique) entre l'étape à forte distribution ramifiée et les transformations en aval pour interrompre la fusion et redistribuer le travail dans le pool de nœuds de calcul. Pour déboguer les problèmes de mémoire associés, consultez Résoudre les erreurs de mémoire insuffisante.
Décalage de clé et touches de raccourci
Les opérations d'agrégation (GroupByKey, CoGroupByKey, Combine.PerKey) regroupent les éléments par clé associée.
- Goulot d'étranglement des clés chaudes : Dataflow achemine tous les éléments ayant la même clé vers un seul thread de travail pour l'agrégation. Si une seule clé contient un pourcentage important de l'ensemble de données, ce nœud de calcul devient un nœud de calcul lent et les nœuds de calcul en amont peuvent subir une contre-pression. Par exemple, une clé
nullpar défaut ou une clé de catégorie extrêmement populaire. - Solution :
- Utilisez des combineurs (
CombineFnouCombine.PerKey) au lieu deGroupByKeydans la mesure du possible, ce qui permet à Dataflow d'effectuer des combinaisons locales partielles avant le shuffle. - Ajoutez un préfixe ou un suffixe entier aléatoire aux touches de raccourci (salage de clé) pour distribuer l'espace de clés entre les nœuds de calcul, suivi d'une agrégation de deuxième niveau pour fusionner les résultats salés.
- Utilisez des combineurs (
Limitation du récepteur en aval
Lorsque vous écrivez la sortie du pipeline dans des services externes tels que des bases de données ou des API tierces, un parallélisme élevé peut saturer le système de destination :
- Limitation : des centaines de threads de nœuds de calcul émettant des appels d'écriture simultanés peuvent entraîner des erreurs de limite de débit, des délais de connexion dépassés ou une dégradation de la base de données.
- Solution :
- Limitez le parallélisme d'écriture en regroupant les éléments avec
GroupByKeyou en utilisant des récepteurs de traitement par lot avec un parallélisme contrôlé. - Implémentez une logique d'intervalle exponentiel entre les tentatives et de nouvelle tentative côté client dans les implémentations de récepteur
DoFn.
- Limitez le parallélisme d'écriture en regroupant les éléments avec
Stratégies d'optimisation
Pour optimiser le parallélisme dans vos jobs Dataflow, envisagez les approches suivantes :
- Empêcher les fusions indésirables avec
Redistribute:Redistribute.arbitrarily(): interrompt la fusion des étapes et redistribue les éléments de manière uniforme entre tous les nœuds de calcul disponibles.Redistribute.byKey(): Rééquilibre les paires clé/valeur entre les threads de nœuds de calcul tout en préservant la localité des clés.- Pour obtenir des exemples d'implémentation, consultez Empêcher la fusion.
- Surveillez les retardataires et les goulots d'étranglement : Utilisez les détails d'exécution de la console Google Cloud pour identifier les étapes avec un nombre élevé de retardataires ou une progression bloquée :
Étapes suivantes
- En savoir plus sur le cycle de vie des pipelines
- Découvrez l'autoscaling horizontal.
- Comprendre le rééquilibrage dynamique du travail
- Consultez les bonnes pratiques pour les pipelines Dataflow.
- Découvrez comment résoudre les erreurs de mémoire insuffisante.