Cette page décrit les caractéristiques de performances, à partir de la version 2.75.0 d'Apache Beam, pour les jobs de streaming Dataflow qui lisent des données à partir d'Apache Kafka et les écrivent dans des tables Apache Iceberg. Elle évalue les différences de performances entre les écritures Apache Iceberg directes et les écritures acheminées via l'API BigQuery gérée, et compare ces résultats aux benchmarks de référence des pipelines Kafka vers BigQuery. Étant donné que les optimisations pour les E/S Apache Iceberg sont en cours, ces métriques de performances sont susceptibles de changer.
Les comparaisons de benchmarks sont disponibles dans trois configurations de mappage sans état principales (ce qui signifie qu'elles lisent à partir de la source, convertissent le message en enregistrement et écrivent dans le récepteur sans suivre l'état ni appliquer de logique métier complexe ; elles sont appelées map_only ou mapping dans les benchmarks) :
- Kafka vers BigQuery (
map_only) (référence des performances de Kafka vers BigQuery) - Kafka vers Iceberg Direct (
map_only,autosharding=false) - Kafka vers Iceberg à l'aide de l'API BigQuery gérée (
map_only)
Ce guide évalue également les modèles de streaming Apache Iceberg directs, tels que le traitement par lot avec état à l'aide de groupbykey, et détaille les considérations en aval essentielles concernant les distributions de taille de fichier, le comportement d'autosharding et la latence des requêtes côté lecture.
Méthodologie de test
Les benchmarks ont été effectués à l'aide des ressources suivantes :
- Cluster Managed Service pour Apache Kafka : le trafic a été généré à l'aide du modèle Dataflow Streaming Data Generator.
- Débit d'entrée : 1 Gbit/s
- Taux de messages : environ 1 000 000 de messages par seconde
- Format des messages : texte JSON avec un schéma fixe (environ 1 Ko par message)
- Partitions : 1 000 partitions Kafka
- Récepteurs de destination
- BigQuery : table standard (non partitionnée) écrite à l'aide de l'API BigQuery Storage Write.
- Apache Iceberg : catalogue basé sur Cloud Storage. Le récepteur direct est partitionné à l'aide de
bucket(id, 64)(mis en bucket dans 64 fragments sur la clé primaire) et utilise le mode de distributionhash.
Une fois l'autoscaling horizontal stabilisé, chaque configuration de pipeline s'est exécutée en état stable pendant 24 heures. Les benchmarks pour chaque cas de pipeline ont été exécutés trois fois, et toutes les valeurs signalées représentent les moyennes calculées sur ces exécutions afin de garantir des métriques de performances fiables et durables.
Performances d'ingestion : charges de travail de mappage
Les pipelines de mappage sans état lisent à partir de la source, convertissent le format du message en enregistrement et écrivent dans le récepteur sans suivre l'état des enregistrements. Les sections suivantes analysent les architectures de référence s'exécutant à 1 Gbit/s.
Configuration du job
| Paramètre | Kafka vers BigQuery (map_only) |
Kafka vers Iceberg Direct (autosharding=false) |
Kafka vers Iceberg à l'aide de l'API BigQuery gérée |
|---|---|---|---|
| Type de machine de nœud de calcul | e2-standard-2 |
e2-standard-4 |
e2-standard-4 |
| vCPU par nœud de calcul | 2 | 4 | 4 |
| RAM par nœud de calcul | 8 Go | 16 Go | 16 Go |
| Streaming Engine | Enabled | Enabled | Enabled |
| Autoscaling horizontal | Enabled | Enabled | Enabled |
| Fréquence de déclenchement | 5 secondes | 60 secondes | 60 secondes |
Débit et utilisation des ressources
L'écriture directe dans des fichiers Parquet physiques dans le stockage d'objets entraîne une surcharge d'E/S plus élevée que l'ingestion en flux continu BigQuery. Par rapport aux écritures Iceberg directes, l'acheminement des écritures via l'API BigQuery gérée améliore l'utilisation du processeur du nœud de calcul (environ 70% contre environ 60%) et réduit légèrement la consommation de Streaming Engine (environ 180 unités de calcul Streaming Engine par heure contre environ 200 unités de calcul Streaming Engine par heure), bien que les exigences globales de calcul du nœud de calcul restent similaires (environ 440 vCPU contre environ 450 vCPU).
| Métrique | Kafka vers BigQuery (map_only) |
Kafka vers Iceberg Direct (autosharding=false) |
Kafka vers Iceberg à l'aide de l'API BigQuery gérée |
|---|---|---|---|
| Débit d'entrée moyen par nœud de calcul | ~15 Mo/s | ~9 Mo/s | ~9 Mo/s |
| Utilisation moyenne du processeur | ~70% | ~60% | ~70% |
| vCPU estimées pour une entrée de 1 Gbit/s | ~126 vCPU | ~450 vCPU | ~440 vCPU |
| Nœuds de calcul estimés pour une entrée de 1 Gbit/s | ~63 nœuds de calcul | ~110 nœuds de calcul | ~110 nœuds de calcul |
| Unités de calcul Streaming Engine estimées par heure pour 1 Gbit/s | ~58 unités de calcul Streaming Engine par heure | ~200 unités de calcul Streaming Engine par heure | ~180 unités de calcul Streaming Engine par heure |
Profil de latence d'écriture
Les écritures Iceberg directes présentent une latence de queue (P99) élevée en raison des contraintes de commit des métadonnées de stockage d'objets. L'utilisation de l'API BigQuery gérée élimine les pics de latence de queue tout en maintenant une latence médiane faible.
| Latence d'écriture de bout en bout | Kafka vers BigQuery | Kafka vers Iceberg Direct (autosharding=false) |
Kafka vers Iceberg à l'aide de l'API BigQuery gérée |
|---|---|---|---|
| P50 (médiane) | ~1 200 ms | ~1 000 ms | ~1 000 ms |
| P95 | ~3 000 ms | ~7 400 ms | ~1 900 ms |
| P99 (queue) | ~5 400 ms | ~14 000 ms | ~2 700 ms |
Considérations sur l'autosharding et choix de conception
Cette section aborde les implications de l'autosharding sur la taille des fichiers et la latence des pipelines lors de l'écriture dans Apache Iceberg.
Pourquoi autosharding=false a été choisi comme référence
Lors des tests initiaux, l'activation de l'autosharding a entraîné la réduction de la taille des fichiers en petits fragments et des fluctuations arbitraires en raison du fractionnement dynamique des fragments déclenché par des pics de charge localisés au niveau des threads, même sous une charge d'entrée agrégée constante.
Pour maintenir des mises en page de fichiers Parquet stables et prévisibles (environ 800 Ko en moyenne) et garantir une référence équitable sans vidages prématurés, autosharding=false a été choisi pour la configuration du récepteur direct.
Que se passe-t-il si vous désactivez l'autosharding au lieu de le laisser activé ?
- Avec
autosharding=false(référence) : vous obtenez des tailles de fichier initiales plus importantes (environ 800 Ko en moyenne) par rapport à l'autosharding. Bien que cela reste faible par rapport aux tailles de fichier Iceberg idéales (128 à 512 Mo), cela nécessite beaucoup moins de compactage en aval. Toutefois, le compromis est une latence de queue d'écriture élevée (P99 atteignant environ 14,0 s) en raison des goulots d'étranglement des métadonnées de stockage d'objets. - Si l'autosharding est activé : Dataflow ajuste dynamiquement les threads d'écriture pour absorber les pics de débit locaux, ce qui réduit la latence de queue d'écriture. Toutefois, cela compromet la couche de stockage en produisant un volume élevé de petits fichiers Parquet fragmentés (environ 100 Ko ou moins). Ces tailles de fichier présentent une variance élevée et fluctuent arbitrairement entre les exécutions (allant d'environ 39 Ko à environ 100 Ko en moyenne), ce qui augmente la nécessité d'une maintenance agressive du compactage en aval.
Réglage et recommandations concernant les partitions
Au cours de notre évaluation, nous avons testé différentes valeurs de partition fixes pour la table de destination afin de trouver un équilibre optimal. Nous avons constaté que l'utilisation de 64 buckets (par exemple, bucket(id, 64)) pour le partitionnement de la table de destination générait les tailles de fichier ciblées tout en maintenant une utilisation et un débit corrects. Cette approche nous a permis d'égaler les avantages en termes de performances de l'autosharding tout en évitant les problèmes de fragmentation de la taille des fichiers arbitraires liés au scaling entièrement dynamique.
Recommandation pour les praticiens : nous encourageons les clients à effectuer des tests préliminaires similaires avec des paramètres de partition ciblés afin de trouver le juste milieu qui maximise le parallélisme des pipelines sans compromettre la taille des fichiers Parquet.
Implications de la lecture en aval : taille des fichiers et compactage
Bien que les métriques côté écriture favorisent l'API BigQuery gérée pour l'ingestion Iceberg, l'efficacité globale du pipeline dépend fortement des performances de lecture en aval :
- Génération de petits fichiers dans l'API BigQuery gérée : l'API BigQuery gérée vide fréquemment les données pour garantir une faible latence d'écriture. Ce comportement entraîne un volume élevé de petits fichiers Parquet écrits dans le catalogue Iceberg cible.
- Impact sur la latence des requêtes de lecture : les moteurs de requête (par exemple, Starburst/Trino, Apache Spark, BigQuery, Dremio) qui lisent des tables contenant des millions de petits fichiers Parquet entraînent une surcharge importante de l'analyse des métadonnées et des pénalités de balayage des partitions.
- Exigences de compactage : pour éviter la dégradation des performances de lecture lors de l'utilisation de l'API BigQuery gérée (ou si l'autosharding est activé sur les écritures directes), exécutez régulièrement des tâches de maintenance de compactage Iceberg (par exemple,
REWRITE DATA FILES). La surcharge de calcul pour le compactage doit être prise en compte dans la conception globale de l'architecture. - Distribution des fichiers d'écriture directe (
autosharding=false) : les écritures Iceberg directes avec un sharding fixe produisent des fichiers Parquet de taille moyenne plus importante (environ 800 Ko), ce qui donne une mise en page moins fragmentée pour un accès immédiat aux requêtes sans demandes de compactage immédiates (bien qu'elles soient toujours inférieures à la plage idéale).
Pipelines Iceberg directs avec état (groupbykey)
Pour évaluer les stratégies de traitement par lot manuel, le regroupement de clés avec état (groupbykey) a été testé par rapport au pipeline de référence Kafka vers Iceberg Direct (map_only, autosharding=false). Les deux configurations écrivent des fichiers Parquet directement dans le stockage d'objets.
Comparaison des benchmarks
| Métrique / Fonctionnalité | Référence du récepteur direct (autosharding=false) |
Récepteur direct avec état (groupbykey) |
Impact sur la performance |
|---|---|---|---|
| vCPU estimées pour 1 Gbit/s | ~450 vCPU | ~520 vCPU | ~+16% de calcul requis |
| Utilisation moyenne du processeur | ~60% | ~50% | ~-17% d'efficacité du nœud de calcul |
| Unités de calcul Streaming Engine estimées par heure pour 1 Gbit/s | ~200 unités de calcul Streaming Engine par heure | ~300 unités de calcul Streaming Engine par heure | ~+50% de charge Streaming Engine |
| Taille moyenne des fichiers | ~800 Ko | ~100 Ko | Génère des lots de fichiers plus petits |
| Latence P50 | ~1 000 ms | ~1 200 ms | ~+20% de médiane plus lente |
| Latence P95 | ~7 400 ms | ~5 500 ms | ~-26% de latence inférieure |
| Latence P99 | ~14 000 ms | ~13 000 ms | Changement marginal de la latence de queue |
Analyse des compromis
- Surcharge de Streaming Engine : l'ajout d'une étape
groupbykeyavec état nécessite que Beam stocke l'état intermédiaire au-delà des limites de la fenêtre. Cela augmente la consommation d'unités de calcul Streaming Engine de ~50% (de ~200 unités de calcul Streaming Engine par heure à ~300 unités de calcul Streaming Engine par heure). - Latence de mise en mémoire tampon : l'agrégation manuelle des clés introduit une mise en mémoire tampon obligatoire des fenêtres, ce qui augmente la latence d'écriture médiane (P50) à ~1 200 ms et la latence P95 à ~5,5 s.
Pipelines inversés : streaming d'Iceberg vers Kafka
Pour évaluer les capacités bidirectionnelles du data lakehouse, des benchmarks ont également été effectués pour le streaming de données en sens inverse : lecture de flux en mode ajout uniquement à partir d'une table Apache Iceberg et publication dans Apache Kafka.
Configuration et efficacité du job
Contrairement aux pipelines d'ingestion qui doivent gérer des écritures de fichiers de stockage d'objets volumineuses ou des goulots d'étranglement de commit des métadonnées, la lecture et le streaming des modifications à partir d'Iceberg fonctionnent avec une efficacité élevée :
| Métrique | Iceberg vers Kafka (ajout uniquement, exactement une fois) |
|---|---|
| Type de machine de nœud de calcul | e2-standard-4 |
| vCPU estimées pour une entrée de 1 Gbit/s | ~30 vCPU |
| Nœuds de calcul estimés pour une entrée de 1 Gbit/s | ~7 nœuds de calcul |
| Unités de calcul Streaming Engine estimées par heure pour 1 Gbit/s | ~0,2 unité de calcul Streaming Engine par heure |
Points clés à retenir concernant les pipelines inversés
- Surcharge de calcul nettement inférieure : la lecture et la projection de flux CDC à partir d'Iceberg nécessitent beaucoup moins de ressources de calcul (environ 30 vCPU contre environ 450 vCPU pour les écritures directes), car elles évitent les tâches lourdes de partitionnement, d'encodage et de commit de grands volumes de fichiers Parquet dans le stockage d'objets.
- Efficacité des ressources : la consommation ou la réplication en aval basée sur les événements à partir d'un format de data lakehouse vers des couches de streaming est très efficace par rapport aux chemins d'ingestion entrants.
Résumé des recommandations architecturales
| Modèle d'architecture | Latence d'écriture P99 | Mise en page du fichier | Considérations sur la lecture en aval |
|---|---|---|---|
Kafka vers BigQuery (map_only) |
~5,4 s | N/A | Optimale (moteur de stockage BigQuery géré) |
| Kafka vers Iceberg à l'aide de l'API BigQuery gérée | ~2,7 s | Fichiers arbitrairement petits | Nécessite un compactage périodique pour les lectures à volume élevé |
Kafka vers Iceberg Direct (autosharding=false) |
~14,0 s | ~800 Ko | Bonne (taille de fichier initiale plus importante, demande de compactage inférieure) |
Kafka vers Iceberg Direct (groupbykey) |
~13,0 s | ~100 Ko | Modérée (surcharge de calcul et d'état plus élevée) |
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.
- Pour 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 le champ Type de job, 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 sur 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 proportionnels 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 des benchmarks.
Supposons, par exemple, que vous exécutiez un pipeline à l'aide de l'architecture Kafka vers Iceberg Direct (autosharding=false) , avec un débit de données d'entrée de 100 Mo/s. En fonction des résultats des benchmarks pour un pipeline de 1 Gbit/s, vous pouvez estimer les besoins en ressources comme suit :
- Facteur de scaling : (100 Mo/s) / (1 024 Mo/s) = ~0,1
- Nœuds de calcul projetés : 110 nœuds de calcul × 0,1 = ~11 nœuds de calcul
- Nombre projeté d'unités de calcul Streaming Engine par heure : 200 × 0,1 = ~20 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 Bonnes pratiques pour l'optimisation des coûts Dataflow.
Exécuter un pipeline de test
Pour déployer un job de streaming Apache Iceberg à l'aide du modèle Flex Dataflow, utilisez la
gcloud dataflow flex-template run
commande.
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--region=REGION \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_To_Iceberg_Yaml \
--enable-streaming-engine \
--parameters ^@^bootstrapServers="KAFKA_BOOTSTRAP_ADDRESS"\
@topic="KAFKA_TOPIC"\
@table="ICEBERG_TABLE_IDENTIFIER"\
@catalogName="CATALOG_NAME"\
@catalogProperties='{"type":"CATALOG_TYPE","warehouse":"gs://BUCKET_NAME/warehouse/"}'\
@triggeringFrequencySeconds=60\
@schema='SCHEMA_DEFINITION'
Remplacez les éléments suivants :
JOB_NAME: nom de votre job DataflowPROJECT_ID: ID de votre Google Cloud projetREGION: Google Cloud région où votre job s'exécute (par exemple,us-central1)KAFKA_BOOTSTRAP_ADDRESS: adresse d'amorçage de votre cluster Apache KafkaKAFKA_TOPIC: nom de votre sujet KafkaICEBERG_TABLE_IDENTIFIER: identifiant de votre table Iceberg cibleCATALOG_NAME: nom de votre catalogue IcebergCATALOG_TYPE: type de catalogue à utiliser (par exemple,hadoopoubigquery)BUCKET_NAME: nom du bucket Cloud Storage pour l'emplacement de votre entrepôtSCHEMA_DEFINITION: définition de schéma pour les données de votre sujet Kafka (par exemple,{"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})