Cette page décrit les caractéristiques de performances des jobs de traitement en flux continu Dataflow qui lisent des données depuis Apache Kafka et les écrivent dans BigQuery. Elle présente les résultats de tests de référence pour les pipelines de carte uniquement, qui effectuent des transformations par message sans suivre l'état ni regrouper les éléments dans le flux.
De nombreuses charges de travail d'intégration de données, y compris l'ETL, la validation de champs et le mappage de schéma, appartiennent à la catégorie "carte uniquement". Si votre pipeline suit ce modèle, vous pouvez utiliser ces benchmarks pour évaluer votre job Dataflow par rapport à une configuration de référence performante.
Méthodologie de test
Les benchmarks ont été effectués à l'aide des ressources suivantes :
Un cluster Managed Service pour Apache Kafka. Les messages ont été générés à l'aide du modèle Streaming Data Generator.
- Débit de messages : environ 1 000 000 de messages par seconde
- Charge d'entrée : 1 Gio/s
- Format des messages : texte JSON généré de manière aléatoire avec un schéma fixe
- Taille des messages : environ 1 Kio par message
- Partitions Kafka: 1 000
Une table BigQuery standard.
Un pipeline de traitement en flux continu Dataflow qui utilise le modèle Apache Kafka vers BigQuery. Ce pipeline effectue l'analyse et le mappage de schéma minimaux requis. Aucune fonction définie par l'utilisateur personnalisée n'a été utilisée.
Une fois le scaling horizontal stabilisé et le pipeline atteint un état stable, les pipelines ont été autorisés à s'exécuter pendant environ une journée, après quoi les résultats ont été collectés et analysés.
Pipeline Dataflow
Ce benchmark utilise un pipeline de carte uniquement qui effectue un mappage et une conversion simples des messages JSON. Le pipeline a été testé en mode "exactement une fois" et en mode "au moins une fois". Le traitement de type "au moins une fois" offre un meilleur débit. Toutefois, il ne doit être utilisé que lorsque les enregistrements en double sont acceptables ou que le récepteur en aval gère la déduplication.
Configuration du job
Le tableau suivant montre comment les jobs Dataflow ont été configurés.
| Paramètre | Valeur |
|---|---|
| Type de machine des nœuds de calcul | e2-standard-2 |
| Processeurs virtuels de nœud de calcul | 2 |
| RAM de la machine de nœud de calcul | 8 Go |
| Disque persistant de la machine de nœud de calcul | Disque persistant standard (HDD), 30 Go |
| Nombre maximal de nœuds de calcul | 120 |
| Streaming Engine | Oui |
| Autoscaling horizontal | Oui |
| Modèle de facturation | Facturation basée sur les ressources |
| L'API Storage Write (gRPC) est-elle activée ? | Oui |
| Flux de l'API Storage Write (gRPC) | 400 |
| Fréquence de déclenchement de l'API Storage Write (gRPC) | 5 secondes |
| Format des messages | JSON |
| Mode d'authentification Kafka |
Identifiants par défaut de l'application (ADC). Pour en savoir plus, consultez la section Types d'authentification pour les courtiers Kafka. |
L'API BigQuery Storage Write (gRPC) est recommandée pour les pipelines de traitement en flux continu. Lorsque vous utilisez le mode "exactement une fois" avec l'API Storage Write (gRPC), vous pouvez ajuster les paramètres suivants :
Nombre de flux d'écriture. Pour garantir un parallélisme de clé suffisant lors de l'étape d'écriture, définissez le nombre de flux de l'API Storage Write (gRPC) sur une valeur supérieure au nombre de processeurs de nœud de calcul, tout en suivant les recommandations de débit par flux.
Fréquence de déclenchement. Une valeur de seconde à un seul chiffre convient aux pipelines à haut débit.
Pour en savoir plus, consultez la section Écrire des données depuis Dataflow vers BigQuery.
Une attention particulière doit également être accordée au nombre de partitions Apache Kafka. Pour garantir un parallélisme de clé suffisant lors de l'étape de lecture, le nombre de partitions doit au moins être égal au nombre total de processeurs virtuels de nœud de calcul. Pour en savoir plus, consultez la section Lire des données depuis Apache Kafka vers Dataflow.
Résultats du benchmark
Cette section décrit les résultats des tests de référence.
Débit et utilisation des ressources
Le tableau suivant présente les résultats des tests pour le débit du pipeline et l'utilisation des ressources.
| Résultat | Exactement une fois | Au moins une fois |
|---|---|---|
| Débit d'entrée par nœud de calcul | Moyenne : 15 Mo/s, n=3 | Moyenne : 18 Mo/s, n=3 |
| Utilisation moyenne du processeur sur tous les nœuds de calcul | Moyenne : 70%, n=3 | Moyenne : 75%, n=3 |
| Nombre de nœuds de calcul | Moyenne : 63, n=3 | Moyenne : 53, n=3 |
| Unités de calcul Streaming Engine par heure | Moyenne : 58, n=3 | Moyenne : 0, n=3 |
L'algorithme d'autoscaling peut affecter le niveau d'utilisation cible du processeur. Pour atteindre une utilisation cible du processeur plus élevée ou plus faible, vous pouvez définir la plage d'autoscaling ou l' indicateur d'utilisation du nœud de calcul. Des cibles d'utilisation plus élevées peuvent entraîner des coûts plus faibles, mais aussi une latence de fin plus mauvaise, en particulier pour les charges variables.
Latence
Le tableau suivant présente les résultats du benchmark pour la latence du pipeline en mode "exactement une fois", à l'exclusion de l'étape d'entrée.
| Latence de bout en bout totale de l'étape, à l'exclusion de l'étape d'entrée | Exactement une fois |
|---|---|
| P50 | Moyenne : 1 200 ms, n=3 |
| P95 | Moyenne : 3 000 ms, n=3 |
| P99 | Moyenne : 5 400 ms, n=3 |
Les tests ont mesuré la latence de bout en bout par étape (la
job/streaming_engine/stage_end_to_end_latencies
métrique) sur trois exécutions de test de longue durée. Cette métrique mesure le temps que Streaming Engine passe à chaque étape du pipeline. Elle englobe toutes les étapes internes du pipeline, telles que :
- Le brassage et la mise en file d'attente des messages pour le traitement
- Le temps de traitement réel, par exemple, la conversion des messages en objets de ligne
- L'écriture de l'état persistant, ainsi que le temps passé en file d'attente pour écrire l'état persistant
En raison d'une limitation de la métrique, la latence de l'étape d'entrée n'est pas signalée. Par conséquent, elle n'est pas incluse dans le total.
Les benchmarks présentés ici représentent une référence. La latence est très sensible à la complexité du pipeline. Les UDF personnalisées, les transformations supplémentaires et la logique de fenêtrage complexe peuvent toutes augmenter la latence.
Coûts estimés
Vous pouvez estimer le coût de référence de votre propre pipeline comparable avec la facturation basée sur les ressources à l'aide du Google Cloud simulateur de coût, comme suit :
- Ouvrez le simulateur de coût.
- Cliquez sur Ajouter à l'estimation.
- Sélectionnez Dataflow.
- Dans Type de service, sélectionnez "Dataflow Classic".
- Sélectionnez Paramètres avancés pour afficher l'ensemble complet des options.
- Choisissez l'emplacement où le job s'exécute.
- Dans Type de mission, sélectionnez "Streaming".
- Sélectionnez Activer Streaming Engine.
- Saisissez les informations concernant les heures d'exécution du job, les nœuds de calcul, les machines de nœud de calcul et le stockage de disque persistant.
- Saisissez le nombre estimé d'unités de calcul Streaming Engine.
L'utilisation des ressources et les coûts sont à peu près linéaires par rapport au débit d'entrée, bien que pour les petits jobs avec seulement quelques nœuds de calcul, le coût total soit dominé par les coûts fixes. Pour commencer, vous pouvez extrapoler le nombre de nœuds de calcul et la consommation de ressources à partir des résultats du benchmark.
Supposons, par exemple, que vous exécutiez un pipeline de carte uniquement en mode "exactement une fois", avec un débit de données d'entrée de 100 Mio/s. En fonction des résultats du benchmark pour un pipeline de 1 Gio/s, vous pouvez estimer les besoins en ressources comme suit :
- Facteur de scaling : (100 Mio/s) / (1 Gio/s) = 0,1
- Nœuds de calcul projetés : 63 nœuds de calcul × 0,1 = 6,3 nœuds de calcul
- Nombre projeté d'unités de calcul Streaming Engine par heure : 58 × 0,1 = 5,8 unités par heure
Cette valeur ne doit être utilisée qu'à titre d'estimation initiale. Le débit et le coût réels peuvent varier considérablement en fonction de facteurs tels que le type de machine, la distribution de la taille des messages, le code utilisateur, le type d'agrégation, le parallélisme des clés et la taille de la fenêtre. Pour en savoir plus, consultez la section Bonnes pratiques pour l'optimisation des coûts Dataflow.
Exécuter un pipeline de test
Cette section présente les
gcloud dataflow flex-template run
commandes qui ont été utilisées pour exécuter le pipeline de carte uniquement.
Mode "exactement une fois"
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_to_BigQuery_Flex \
--enable-streaming-engine \
--parameters \
readBootstrapServerAndTopic="KAFKA_BOOTSTRAP_ADDRESS;KAFKA_TOPIC",\
kafkaReadAuthenticationMode=APPLICATION_DEFAULT_CREDENTIALS,\
messageFormat=JSON,\
writeMode=SINGLE_TABLE_NAME,\
outputTableSpec="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME",\
useBigQueryDLQ=true,\
outputDeadletterTable="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME_dlq",\
numStorageWriteApiStreams=400
Mode "au moins une fois"
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_to_BigQuery_Flex \
--enable-streaming-engine \
--additional-experiments=streaming_mode_at_least_once \
--parameters \
readBootstrapServerAndTopic="KAFKA_BOOTSTRAP_ADDRESS;KAFKA_TOPIC",\
kafkaReadAuthenticationMode=APPLICATION_DEFAULT_CREDENTIALS,\
messageFormat=JSON,\
writeMode=SINGLE_TABLE_NAME,\
outputTableSpec="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME",\
useBigQueryDLQ=true,\
outputDeadletterTable="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME_dlq",\
numStorageWriteApiStreams=400,\
useStorageWriteApiAtLeastOnce=true
Remplacez les éléments suivants :
JOB_NAME: nom du job DataflowPROJECT_ID: ID du projetKAFKA_BOOTSTRAP_ADDRESS: adresse d'amorçage du cluster Apache KafkaKAFKA_TOPIC: nom du sujet KafkaBQ_DATASET: nom de l'ensemble de données BigQueryBQ_TABLE_NAME: nom de la table BigQuery
Générer des données test
Pour générer des données de test, exécutez la commande suivante pour exécuter le modèle Streaming Data Generator :
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Streaming_Data_Generator \
--max-workers=140 \
--parameters \
schemaLocation=SCHEMA_LOCATION,\
qps=1000000,\
sinkType=KAFKA,\
bootstrapServer=KAFKA_BOOTSTRAP_ADDRESS,\
kafkaTopic=KAFKA_TOPIC,\
outputType=JSON
Remplacez les éléments suivants :
JOB_NAME: nom du job DataflowPROJECT_ID: ID du projetSCHEMA_LOCATION: chemin d'accès à un fichier de schéma dans Cloud StorageKAFKA_BOOTSTRAP_ADDRESS: adresse d'amorçage du cluster Apache KafkaKAFKA_TOPIC: nom du sujet Kafka
Le modèle Streaming Data Generator utilise un fichier JSON Data Generator pour définir le schéma des messages. Les tests de référence ont utilisé un schéma de message semblable à celui-ci :
{ "logStreamId": "{{integer(1000001,2000000)}}", "message": "{{alphaNumeric(962)}}" }
Étapes suivantes
- Utiliser l'interface de surveillance des jobs Dataflow
- Bonnes pratiques pour l'optimisation des coûts Dataflow
- Résoudre les problèmes liés aux jobs de traitement en flux continu lents ou bloqués
- Lire des données depuis Apache Kafka vers Dataflow
- Écrire des données depuis Dataflow vers BigQuery