Résoudre les problèmes de goulots d'étranglement dans Dataflow

Un goulot d'étranglement se produit lorsqu'une étape, une phase ou un nœud de calcul ralentit le job dans son ensemble. Les goulots d'étranglement peuvent entraîner l'inactivité des nœuds de calcul et une latence accrue.

Si Dataflow détecte un goulot d'étranglement, le graphique du job affiche une alerte, et le panneau Step Info (Informations sur l'étape) indique le type de goulot d'étranglement et la cause, si elle est connue. Dataflow exporte également les informations de détection des goulots d'étranglement vers une métrique Monitoring, qui présente les données sous forme de série temporelle. Vous pouvez ainsi afficher les goulots d'étranglement au fil du temps ou dans le passé.

Comprendre les goulots d'étranglement

Lorsque Dataflow exécute un pipeline de traitement en flux continu, le job se compose d'une série de composants, tels que des brassages en flux continu, des threads de traitement de fonctions définies par l'utilisateur (DoFn) et des points de contrôle d'état persistants. Pour faciliter le flux de données, Dataflow utilise des files d'attente pour connecter ces composants. Les données sont transmises de l'amont vers l'aval.

Dans de nombreux pipelines, la capacité de débit globale est limitée par un seul composant, ce qui crée un goulot d'étranglement dans le pipeline. Le débit auquel les données peuvent transiter par un goulot d'étranglement limite la vitesse à laquelle le pipeline peut accepter et traiter les données d'entrée.

Prenons l'exemple d'un pipeline dans lequel le traitement DoFn a lieu en aval d'un brassage en flux continu. Une file d'attente entre eux met en mémoire tampon les données brassées, mais non traitées. Si le traitement DoFn ne peut pas consommer les données aussi rapidement que le brassage en flux continu les produit, la file d'attente augmente. Un goulot d'étranglement prolongé peut entraîner la saturation de la file d'attente. À ce stade, le brassage est mis en pause et le backlog se propage en amont. Les files d'attente en amont accumulent également des backlogs, ce qui finit par entraîner un ralentissement qui s'étend à la source de données. Cela signifie que l'ensemble du pipeline ne peut pas suivre le rythme de l'entrée.

Lorsqu'un goulot d'étranglement se produit, une partie importante du pipeline peut sembler défectueuse, même si un seul point du pipeline est à l'origine du backlog. Ce comportement peut rendre le débogage des goulots d'étranglement difficile. L'objectif de la détection des goulots d'étranglement est d'identifier l'emplacement et la cause exacts, en éliminant les conjectures, afin que vous puissiez résoudre la cause première.

Dataflow détecte un goulot d'étranglement lorsqu'un délai dépasse le seuil de cinq minutes. Si le délai ne dépasse pas ce seuil, Dataflow ne détecte pas de goulot d'étranglement.

La détection des goulots d'étranglement ne nécessite pas toujours d'action de votre part et dépend de votre cas d'utilisation. Un pipeline peut fonctionner normalement avec des délais temporaires de plus de cinq minutes. Si cela est acceptable pour votre cas d'utilisation, vous n'aurez peut-être pas besoin de résoudre les goulots d'étranglement indiqués.

Types de goulots d'étranglement

Lorsque Dataflow détecte un goulot d'étranglement, l'interface de surveillance indique la gravité du problème. Les goulots d'étranglement se répartissent dans les catégories suivantes :

Le traitement est bloqué et ne progresse pas.
La progression du pipeline est complètement interrompue à cette étape.
Le traitement est en cours, mais il est en retard.
Le pipeline ne peut pas traiter les données entrantes aussi rapidement qu'elles arrivent. Le backlog augmente donc.
Le traitement est en cours, mais le backlog est stable.
Le pipeline progresse et le taux de traitement est comparable au taux d'entrée. Le traitement est suffisamment rapide pour que le backlog n'augmente pas, mais le backlog accumulé ne diminue pas non plus de manière significative.
Le traitement est en cours et rattrape le retard accumulé.
Le backlog diminue, mais le goulot d'étranglement actuel empêche le pipeline de rattraper son retard plus rapidement. Si vous démarrez un pipeline avec un backlog, cet état peut être normal et ne nécessiter aucune intervention. Surveillez la progression pour voir si le backlog continue de diminuer.

Causes des goulots d'étranglement

Cette section répertorie les causes des goulots d'étranglement qui peuvent être détectées. Utilisez ces informations pour résoudre le problème. Dans certains cas, plusieurs causes peuvent être présentes et liées. Par exemple, si les nœuds de calcul sont sous-provisionnés, l'utilisation des vCPU peut être élevée. Une utilisation élevée des vCPU peut ralentir les opérations, ce qui peut entraîner un délai de mise en file d'attente plus long. L'analyse des causes probables peut afficher toutes ces causes comme étant à l'origine du goulot d'étranglement.

Opérations de traitement de longue durée

Certaines opérations de ce calcul ont une longue durée de traitement. Cela se produit chaque fois qu'un bundle d'entrée est envoyé au nœud de calcul exécutant le DoFn et qu'un délai important s'est écoulé sans que des résultats soient disponibles.

Cela est le plus souvent dû à une seule opération de longue durée dans le code utilisateur. D'autres problèmes peuvent se manifester sous la forme d'opérations de traitement de longue durée. Par exemple, les erreurs générées et relancées dans le DoFn, les tentatives pendant de longues périodes ou les plantages du harnais du nœud de calcul en raison de facteurs tels que les erreurs de mémoire saturée peuvent tous entraîner ces longues durées de traitement.

Si le calcul concerné se trouve dans le code utilisateur, recherchez des moyens d'optimiser le code ou de limiter la durée d'exécution. Pour faciliter le débogage, les journaux des nœuds de calcul affichent des traces de pile pour toutes les opérations bloquées pendant plus de 5 minutes.

Commit de clé trop volumineux

Consultez la section Exception liée à un commit de clé trop volumineux.

Temps de traitement long pour toutes les opérations

Les opérations de ce calcul prennent systématiquement beaucoup de temps, ce qui suggère un problème dans le DoFn fourni par l'utilisateur.

Cette cause est différente des opérations de traitement de longue durée. Alors que cette cause affecte certaines opérations, celle-ci indique que toutes les opérations de ce calcul sont concernées.

Recherchez dans les journaux des nœuds de calcul les erreurs, les exceptions ou les traces de pile indiquant des threads lents ou bloqués. Si vous utilisez le SDK Apache Beam pour Python et que les opérations sont intrinsèquement longues par conception (par exemple, appels d'API externes lents ou E/S à latence élevée), envisagez d'utiliser un DoFn asynchrone. Cette fonctionnalité peut améliorer le débit en empêchant le traitement d'être bloqué sur ces tâches de longue durée.

Lecture lente de l'état persistant

Le calcul passe beaucoup de temps à lire l'état persistant lors de l'exécution du DoFn. Cela peut être dû à un état persistant excessivement volumineux ou à un nombre trop important de lectures. Envisagez de réduire la taille de l'état persistant ou la fréquence des lectures. Optimisez les modèles d'accès à l'état en combinant plusieurs opérations de lecture d'état ou en utilisant l'état de mappage au lieu de l'état de valeur, le cas échéant. Il peut également s'agir d'un problème temporaire dû à la lenteur de l'état persistant sous-jacent.

Écriture lente de l'état persistant

Le calcul passe beaucoup de temps à écrire l'état persistant lors du commit des résultats du traitement. Cela peut être dû à un état persistant excessivement volumineux. Envisagez de réduire la taille de l'état persistant. Optimisez les modèles d'accès à l'état en combinant plusieurs opérations d'écriture d'état, le cas échéant. Il peut également s'agir d'un problème temporaire dû à la lenteur de l'état persistant sous-jacent.

Commit refusé

Le traitement des données ne peut pas être validé dans l'état persistant, car il n'est pas valide. Cela est généralement dû au dépassement d'une des limites opérationnelles limites. Pour en savoir plus, consultez les journaux ou contactez l'assistance.

Nombre insuffisant de partitions sources Apache Kafka

Le calcul de la source Apache Kafka ne comporte pas suffisamment de partitions. Pour résoudre ce problème, procédez comme suit :

  • Augmentez le nombre de partitions Kafka.
  • Incluez la redistribution à l'aide de .withRedistribute() lors de la configuration de la lecture des E/S Kafka pour paralléliser les données plus efficacement. Incluez .withRedistributeNumKeys(N)N > partitions pour fournir une limite supérieure au nombre total de clés. Un nombre limité de clés permet d'améliorer l'efficacité grâce au regroupement des enregistrements.
  • Pour minimiser le coût du brassage de redistribution, utilisez .withOffsetDeduplication(). Ce mode minimise la quantité de données qui doivent être conservées dans le cadre du brassage, tout en fournissant un traitement exactement une fois.

Pour en savoir plus, consultez la section Parallélisme de la page Lire des données depuis Apache Kafka vers Dataflow.

Volume important d'état persistant de la source Apache Kafka

Le calcul de la source Apache Kafka redistribue un volume élevé de données qui peuvent entraîner une latence et des coûts élevés. Pour résoudre ce problème, procédez comme suit :

  • Si un traitement exactement une fois est requis pour le pipeline, minimisez le coût du brassage de redistribution en utilisant le mode de déduplication des décalages. Ce mode minimise la quantité de données qui doivent être conservées dans le cadre du brassage, tout en fournissant un traitement exactement une fois.
  • Si un traitement au moins une fois est suffisant pour le pipeline, la configuration permettant les doublons peut être activée.

Pour en savoir plus, consultez la page Lire des données depuis Apache Kafka vers Dataflow.

Parallélisme de la source insuffisant

Un calcul de source présente un parallélisme insuffisant. Si possible, augmentez la parallélisation dans la source. Si vous ne pouvez pas augmenter la parallélisation et que le job utilise le mode au moins une fois, essayez d'ajouter une transformation Redistribute au pipeline.

Clés surutilisées ou parallélisme de clés insuffisant

Le job comporte des clés surutilisées ou un parallélisme de clés insuffisant.

Pour chaque clé de partitionnement, Dataflow traite les messages de manière séquentielle. Pendant que Dataflow traite un lot de messages pour une clé donnée, les autres messages entrants pour cette clé sont mis en file d'attente jusqu'à ce que le lot actuel soit terminé.

Si Dataflow ne peut pas traiter suffisamment de clés distinctes en parallèle, cela peut entraîner un goulot d'étranglement. Par exemple, les données peuvent avoir trop peu de clés distinctes, ou certaines clés peuvent être surreprésentées dans les données ("clés surutilisées"). Pour résoudre ce problème, modifiez la logique de votre pipeline afin de redistribuer les données. Par exemple, ajoutez un suffixe aléatoire aux clés de regroupement pour diviser les clés surutilisées et agrégez les résultats lors d'une étape ultérieure. Pour en savoir plus, consultez la section Résoudre les problèmes liés aux clés surutilisées.

vCPU sous-provisionnés

Le job ne dispose pas de suffisamment de vCPU de nœud de calcul. Cette situation se produit lorsque le job est déjà mis à l'échelle au maximum, que l'utilisation des vCPU est élevée et qu'il existe toujours un backlog. Vous devrez peut-être augmenter le nombre maximal de nœuds de calcul provisionnés pour ce job. Par exemple, vous pouvez augmenter ce nombre en mettant à jour la plage d'autoscaling . Vous pouvez également rechercher des moyens de réduire l'utilisation des vCPU en modifiant le code du pipeline ou la charge de travail. Vous pouvez utiliser Cloud Profiler pour rechercher des opportunités d'optimisation.

Utilisation élevée des vCPU, en attente d'un scaling à la hausse

Le job présente une utilisation élevée des vCPU, mais il est possible d'effectuer un scaling à la hausse. Cette condition est probablement temporaire jusqu'à ce que le scaling à la hausse puisse avoir lieu. Vous pouvez surveiller l'autoscaling pour voir les décisions d'autoscaling. Si cette condition persiste pendant une longue période ou se produit fréquemment, vous devrez peut-être modifier la configuration de l'autoscaling en définissant un indicateur d'utilisation différent pour permettre au job d'effectuer un scaling à la hausse de manière plus proactive.

Charge vCPU déséquilibrée créant des goulots d'étranglement sur certains nœuds de calcul aberrants

Le job dispose de suffisamment de vCPU de nœud de calcul, mais certains nœuds de calcul affichent une utilisation très élevée des vCPU. Cela est souvent dû à une répartition inégale du travail. Les causes potentielles incluent des partitions sources chargées de manière inégale ou des clés surutilisées.

Pour résoudre ce problème, procédez comme suit :

  • Déterminez la cause du chargement inégal et essayez de la corriger. Par exemple, assurez-vous que les partitions sources sont réparties de manière égale.
  • Si la correction de la charge inégale n'est pas possible, envisagez de modifier la forme de la VM de nœud de calcul pour augmenter le nombre de vCPU par nœud de calcul afin de réduire l'utilisation maximale. Pour en savoir plus sur la configuration des vCPU par nœud de calcul, consultez la section Configurer les VM de nœud de calcul Dataflow.
Problème de communication avec les nœuds de calcul

Dataflow ne peut pas communiquer avec toutes les VM de nœud de calcul. Vérifiez l'état des VM de nœud de calcul du job. Les causes possibles sont les suivantes :

  • Un problème est survenu lors du provisionnement des VM de nœud de calcul.
  • Le pool de VM de nœud de calcul est supprimé pendant l'exécution du job.
  • Problèmes de Mise en réseau.
La source Pub/Sub présente des erreurs de récupération.

Des erreurs se produisent lors de la récupération à partir de la source Pub/Sub. Vérifiez que le sujet et les abonnements requis existent, puis vérifiez le quota et la configuration. Vous pouvez également rechercher des erreurs dans les journaux.

La source Pub/Sub présente un parallélisme insuffisant

Le calcul de la source Pub/Sub ne comporte pas suffisamment de clés Pub/Sub. Pour augmenter le nombre de clés, définissez l'option de service num_pubsub_keys. Pour en savoir plus, consultez la section Parallélisme de la source Pub/Sub.

La source Pub/Sub est limitée pour une raison inconnue

Le calcul de la source Pub/Sub est limité lors de la lecture à partir de Pub/Sub, pour une raison inconnue. Ce problème peut être temporaire. Recherchez les problèmes de configuration Pub/Sub, les autorisations IAM manquantes ou les limites de quota. Toutefois, si aucune des zones précédentes n'est la cause première et que le problème persiste, contactez l'assistance.

La publication au niveau du récepteur Pub/Sub est lente ou bloquée

Le calcul du récepteur Pub/Sub est lent ou bloqué. Ce problème peut être dû à un problème de configuration ou à une limite de quota.

Temps de file d'attente de travail élevé

L'âge de travail éligible le plus ancien est élevé, en raison du grand nombre de clés et du taux de traitement des clés. Dans ce cas, chaque opération n'est peut-être pas anormalement longue, mais le délai de mise en file d'attente global est élevé.

Dataflow utilise un seul thread de traitement par clé de partitionnement, et le nombre de threads de traitement est limité. Le délai de mise en file d'attente est approximativement égal au rapport entre les clés et les threads, multiplié par la latence sur le thread pour chaque bundle de traitement d'une clé :

(key count / total harness threads) * latency per bundle

Vous pouvez essayer les solutions suivantes :

  • Augmenter le nombre de nœuds de calcul. Consultez la section Autoscaling de flux.
  • Augmentez le nombre de threads de harnais de nœud de calcul. Définissez l'option de pipeline numberOfWorkerHarnessThreads / number_of_worker_harness_threads .
  • Réduisez le nombre de clés.
  • Réduisez la latence de l'opération.
Problème temporaire avec le backend Streaming Engine

Un problème de configuration ou opérationnel est survenu avec le backend Streaming Engine. Ce problème peut être temporaire. Si ce problème persiste, contactez l'assistance.

Cause indéterminée

La cause du backlog ne peut pas être déterminée avec certitude. Ce problème peut être temporaire. Si ce problème persiste, contactez l'assistance.

Métriques des goulots d'étranglement

Les métriques de job suivantes fournissent des informations sur les goulots d'étranglement :

Étape suivante