Caractéristiques de performances des pipelines Kafka vers Iceberg

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) :

  1. Kafka vers BigQuery (map_only) (référence des performances de Kafka vers BigQuery)
  2. Kafka vers Iceberg Direct (map_only, autosharding=false)
  3. 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 distribution hash.

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

  1. Surcharge de Streaming Engine : l'ajout d'une étape groupbykey avec é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).
  2. 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 :

  1. Ouvrez le simulateur de coût.
  2. Cliquez sur Ajouter à l'estimation.
  3. Sélectionnez Dataflow.
  4. Pour Type de service, sélectionnez "Dataflow Classic".
  5. Sélectionnez Paramètres avancés pour afficher l'ensemble complet des options.
  6. Choisissez l'emplacement où le job s'exécute.
  7. Dans le champ Type de job, sélectionnez "Streaming".
  8. Sélectionnez Activer Streaming Engine.
  9. 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.
  10. 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 Dataflow
  • PROJECT_ID: ID de votre Google Cloud projet
  • REGION: 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 Kafka
  • KAFKA_TOPIC : nom de votre sujet Kafka
  • ICEBERG_TABLE_IDENTIFIER : identifiant de votre table Iceberg cible
  • CATALOG_NAME : nom de votre catalogue Iceberg
  • CATALOG_TYPE : type de catalogue à utiliser (par exemple, hadoop ou bigquery)
  • BUCKET_NAME : nom du bucket Cloud Storage pour l'emplacement de votre entrepôt
  • SCHEMA_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"}]})