Leistungsmerkmale von Kafka-zu-Iceberg-Pipelines

Auf dieser Seite werden die Leistungsmerkmale von Dataflow-Streamingjobs beschrieben, die aus Apache Kafka lesen und in Apache Iceberg-Tabellen schreiben. Die Daten stammen aus Apache Beam Version 2.75.0. Dabei werden die Leistungsunterschiede zwischen direkten Apache Iceberg-Schreibvorgängen und Schreibvorgängen, die über die verwaltete BigQuery API weitergeleitet werden, bewertet und diese Ergebnisse mit den Referenz-Benchmarks von Kafka zu BigQuery-Pipelines verglichen. Da die Optimierungen für Apache Iceberg-E/A-Vorgänge noch nicht abgeschlossen sind, können sich diese Leistungsmesswerte ändern.

Benchmark-Vergleiche sind für drei primäre zustandslose Zuordnungskonfigurationen verfügbar. Das bedeutet, dass Daten aus der Quelle gelesen, die Nachricht in einen Datensatz konvertiert und in die Senke geschrieben werden, ohne den Status zu verfolgen oder komplexe Geschäftslogik anzuwenden. In den Benchmarks werden diese Konfigurationen als map_only oder mapping bezeichnet.

  1. Kafka zu BigQuery (map_only) (Referenz für die Leistung von Kafka zu BigQuery)
  2. Kafka zu Iceberg direkt (map_only, autosharding=false)
  3. Kafka zu Iceberg mit der verwalteten BigQuery API (map_only)

Außerdem werden in diesem Leitfaden direkte Apache Iceberg-Streamingmuster wie zustandsorientierte Batchverarbeitung mit groupbykey bewertet und wichtige nachgelagerte Aspekte in Bezug auf Dateigrößenverteilungen, Autosharding-Verhalten und Abfragelatenz auf der Leseseite beschrieben.

Testmethodik

Die Benchmarks wurden mit den folgenden Ressourcen durchgeführt:

  • Managed Service for Apache Kafka-Cluster:Der Traffic wurde mit der Dataflow-Vorlage „Streaming Data Generator“ generiert.
    • Eingabedurchsatz:1 GB/s
    • Nachrichtenrate:ca. 1.000.000 Nachrichten pro Sekunde
    • Nachrichtenformat:JSON-Text mit einem festen Schema (ca. 1 KB pro Nachricht)
    • Partitionen:1.000 Kafka-Partitionen
  • Zielsenken:
    • BigQuery:Standardtabelle (nicht partitioniert), die mit der BigQuery Storage Write API geschrieben wurde.
    • Apache Iceberg:Katalog, der von Cloud Storage unterstützt wird. Die direkte Senke wird mit bucket(id, 64) partitioniert (in 64 Shards auf dem Primärschlüssel aufgeteilt) und verwendet den Verteilungsmodus hash.

Nachdem sich das horizontale Autoscaling stabilisiert hatte, wurde jede Pipelinekonfiguration 24 Stunden lang im stabilen Zustand ausgeführt. Die Benchmarks für jeden Pipelinefall wurden dreimal separat ausgeführt. Alle angegebenen Werte stellen die berechneten Durchschnittswerte aus diesen Ausführungen dar, um nachhaltige und zuverlässige Leistungsmesswerte zu gewährleisten.

Aufnahmeleistung: Zuordnungsarbeitslasten

Zustandslose Zuordnungspipelines lesen Daten aus der Quelle, konvertieren das Nachrichtenformat in einen Datensatz und schreiben Daten in die Senke, ohne den Status über Datensätze hinweg zu verfolgen. In den folgenden Abschnitten werden Referenzarchitekturen analysiert, die mit 1 GB/s ausgeführt werden.

Jobkonfiguration

Einstellung Kafka zu BigQuery (map_only) Kafka zu Iceberg direkt (autosharding=false) Kafka zu Iceberg mit der verwalteten BigQuery API
Worker-Maschinentyp e2-standard-2 e2-standard-4 e2-standard-4
vCPUs pro Worker 2 4 4
RAM pro Worker 8 GB 16 GB 16 GB
Streaming Engine Aktiviert Aktiviert Aktiviert
Horizontales Autoscaling Aktiviert Aktiviert Aktiviert
Auslösehäufigkeit 5 Sekunden 60 Sekunden 60 Sekunden

Durchsatz und Ressourcennutzung

Das direkte Schreiben in physische Parquet-Dateien im Objektspeicher verursacht einen höheren E/A-Overhead als die BigQuery-Streamingaufnahme. Im Vergleich zu direkten Iceberg-Schreibvorgängen verbessert die Weiterleitung von Schreibvorgängen über die verwaltete BigQuery API die CPU-Auslastung der Worker (ca. 70% gegenüber ca. 60%) und reduziert den Streaming Engine-Verbrauch geringfügig (ca. 180 SECU/h gegenüber ca. 200 SECU/h). Die gesamten Rechenanforderungen der Worker bleiben jedoch ähnlich (ca. 440 vCPUs gegenüber ca. 450 vCPUs).

Messwert Kafka zu BigQuery (map_only) Kafka zu Iceberg direkt (autosharding=false) Kafka zu Iceberg mit der verwalteten BigQuery API
Durchschnittlicher Eingabedurchsatz pro Worker ca. 15 MB/s ca. 9 MB/s ca. 9 MB/s
Durchschnittliche CPU-Auslastung ca. 70% ca. 60% ca. 70%
Geschätzte vCPUs für 1 GB/s-Eingabe ca. 126 vCPUs ca. 450 vCPUs ca. 440 vCPUs
Geschätzte Worker für 1 GB/s-Eingabe ca. 63 Worker ca. 110 Worker ca. 110 Worker
Geschätzte SECU pro Stunde für 1 GB/s ca. 58 SECU/h ca. 200 SECU/h ca. 180 SECU/h

Profil der Schreiblatenz

Direkte Iceberg-Schreibvorgänge weisen aufgrund von Einschränkungen bei der Commit-Ausführung von Objektspeicher-Metadaten eine hohe Extremwertlatenz (P99) auf. Durch die Verwendung der verwalteten BigQuery API werden Spitzen bei der Extremwertlatenz vermieden und gleichzeitig eine niedrige Medianlatenz beibehalten.

End-to-End-Schreiblatenz Kafka zu BigQuery Kafka zu Iceberg direkt (autosharding=false) Kafka zu Iceberg mit der verwalteten BigQuery API
P50 (Median) ca. 1.200 ms ca. 1.000 ms ca. 1.000 ms
P95 ca. 3.000 ms ca. 7.400 ms ca. 1.900 ms
P99 (Extremwert) ca. 5.400 ms ca. 14.000 ms ca. 2.700 ms

Überlegungen zu Autosharding und Designentscheidungen

In diesem Abschnitt werden die Auswirkungen von Autosharding auf Dateigrößen und die Pipeline-Latenz beim Schreiben in Apache Iceberg erläutert.

Warum autosharding=false als Referenz gewählt wurde

In ersten Tests führte die Aktivierung von Autosharding dazu, dass die Dateigrößen in winzige Blöcke zerfielen und willkürlich schwankten. Das lag an der dynamischen Aufteilung von Shards, die durch lokale Lastspitzen auf Thread-Ebene ausgelöst wurde – selbst bei einer konstanten aggregierten Eingabelast.

Um stabile, vorhersehbare Parquet-Dateilayouts (durchschnittlich ca. 800 KB) beizubehalten und eine faire Referenz ohne vorzeitige Flushes zu gewährleisten, wurde autosharding=false für die Konfiguration der direkten Senke ausgewählt.

Was passiert, wenn Sie Autosharding deaktivieren oder aktiviert lassen?

  • Mit autosharding=false (Referenz): Sie erzielen größere anfängliche Dateigrößen (durchschnittlich ca. 800 KB) als mit Autosharding. Das ist zwar immer noch klein im Vergleich zu idealen Iceberg-Dateigrößen (128–512 MB), erfordert aber deutlich weniger nachgelagerte Komprimierung. Der Nachteil ist jedoch eine hohe Extremwertlatenz beim Schreiben (P99 erreicht ca.14, 0 s) aufgrund von Engpässen bei den Objektspeicher-Metadaten.
  • Wenn Autosharding aktiviert ist:Dataflow skaliert Writer-Threads dynamisch, um lokale Durchsatzspitzen aufzufangen. Dadurch wird die Extremwertlatenz beim Schreiben reduziert. Dies beeinträchtigt jedoch die Speicherebene, da eine große Anzahl kleiner, fragmentierter Parquet-Dateien (ca. 100 KB oder kleiner) erstellt wird. Diese Dateigrößen weisen eine hohe Varianz auf und schwanken willkürlich zwischen den Ausführungen (durchschnittlich ca. 39 KB bis ca. 100 KB). Dadurch steigt der Bedarf an aggressiver nachgelagerter Komprimierung.

Partitionierung optimieren und Empfehlungen

Während unserer Bewertung haben wir mit verschiedenen festen Partitionswerten für die Zieltabelle experimentiert, um ein optimales Gleichgewicht zu finden. Wir haben festgestellt, dass die Verwendung von 64 Buckets (Beispiel: bucket(id, 64)) für die Partitionierung der Zieltabelle die gewünschten Dateigrößen bei angemessener Auslastung und angemessenem Durchsatz lieferte. Mit diesem Ansatz konnten wir die Leistungsvorteile von Autosharding nutzen und gleichzeitig die Probleme mit der willkürlichen Fragmentierung der Dateigröße vermeiden, die mit einer vollständig dynamischen Skalierung verbunden sind.

Empfehlung für Praktiker:Kunden sollten ähnliche Vortests mit gezielten Partitionseinstellungen durchführen, um den optimalen Punkt zu finden, an dem die Pipeline-Parallelität maximiert wird, ohne die Parquet-Dateigrößen zu beeinträchtigen.

Auswirkungen auf das Lesen nachgelagerter Daten: Dateigrößen und Komprimierung

Während die Messwerte auf der Schreibseite die verwaltete BigQuery API für die Iceberg-Aufnahme bevorzugen, hängt die Gesamteffizienz der Pipeline stark von der Leseleistung nachgelagerter Daten ab:

  • Erstellung kleiner Dateien in der verwalteten BigQuery API:Die verwaltete BigQuery API führt häufig Flushes durch, um eine niedrige Schreiblatenz zu gewährleisten. Dieses Verhalten führt zu einer großen Anzahl kleiner Parquet-Dateien, die in den Ziel-Iceberg-Katalog geschrieben werden.
  • Auswirkungen auf die Latenz von Leseabfragen:Abfrage-Engines (z. B. Starburst/Trino, Apache Spark, BigQuery, Dremio), die Tabellen mit Millionen kleiner Parquet-Dateien lesen, verursachen einen hohen Overhead beim Parsen von Metadaten und Strafen beim Scannen von Partitionen.
  • Anforderungen an die Komprimierung:Um eine Beeinträchtigung der Leseleistung bei Verwendung der verwalteten BigQuery API zu vermeiden (oder wenn Autosharding für direkte Schreibvorgänge aktiviert ist), führen Sie regelmäßig Iceberg-Komprimierungsjobs aus (z. B. REWRITE DATA FILES). Der Rechenaufwand für die Komprimierung sollte in das Gesamtdesign der Architektur einbezogen werden.
  • Dateiverteilung bei direkten Schreibvorgängen (autosharding=false):Direkte Iceberg-Schreibvorgänge mit fester Shard-Anzahl erzeugen größere durchschnittliche Parquet-Dateien (ca. 800 KB), was zu einem weniger fragmentierten Layout für den sofortigen Abfragezugriff führt, ohne dass eine sofortige Komprimierung erforderlich ist (obwohl die Dateigröße immer noch unter dem idealen Bereich liegt).

Zustandsorientierte direkte Iceberg-Pipelines (groupbykey)

Um manuelle Batchverarbeitungsstrategien zu bewerten, wurde die zustandsorientierte Gruppierung von Schlüsseln (groupbykey) mit der Referenzpipeline Kafka zu Iceberg direkt (map_only, autosharding=false) verglichen. Bei beiden Konfigurationen werden Parquet-Dateien direkt in den Objektspeicher geschrieben.

Benchmark-Vergleich

Messwert / Funktion Referenz für die direkte Senke (autosharding=false) Zustandsorientierte direkte Senke (groupbykey) Auswirkungen auf die Leistung
Geschätzte vCPUs für 1 GB/s ca. 450 vCPUs ca. 520 vCPUs ca. 16% mehr Rechenleistung erforderlich
Durchschnittliche CPU-Auslastung ca. 60% ca. 50% ca. 17% geringere Worker-Effizienz
Geschätzte SECU/h für 1 GB/s ca. 200 SECU/h ca. 300 SECU/h ca. 50% höhere Streaming Engine-Last
Durchschnittliche Dateigröße ca. 800 KB ca. 100 KB Erstellt kleinere Dateibatches
P50-Latenz ca. 1.000 ms ca. 1.200 ms ca. 20% langsamerer Median
P95-Latenz ca. 7.400 ms ca. 5.500 ms ca. 26% geringere Latenz
P99-Latenz ca. 14.000 ms ca. 13.000 ms Geringfügige Änderung der Extremwertlatenz

Kompromissanalyse

  1. Streaming Engine-Aufwand:Wenn Sie einen zustandsorientierten groupbykey-Schritt hinzufügen, muss Beam den Zwischenstatus über Fenstergrenzen hinweg speichern. Dadurch erhöht sich der Verbrauch von Streaming Engine-Recheneinheiten um ca. 50% (von ca. 200 SECU/h auf ca. 300 SECU/h).
  2. Pufferlatenz:Die manuelle Schlüsselaggregation führt zu einer obligatorischen Fensterpufferung, wodurch sich die Medianlatenz beim Schreiben (P50) auf ca. 1.200 ms und die P95-Latenz auf ca. 5,5 s erhöht.

Umgekehrte Pipelines: Streaming von Iceberg zu Kafka

Um die bidirektionalen Lakehouse-Funktionen zu bewerten, wurden auch Benchmarks für das Streaming von Daten in umgekehrter Richtung durchgeführt. Dabei wurden Streams, an die nur Daten angehängt werden können, aus einer Apache Iceberg-Tabelle gelesen und wieder in Apache Kafka veröffentlicht.

Jobkonfiguration und Effizienz

Im Gegensatz zu Aufnahmepipelines, bei denen viele Dateischreibvorgänge im Objektspeicher oder Engpässe bei der Commit-Ausführung von Metadaten auftreten können, ist das Lesen und Streamen von Änderungen aus Iceberg sehr effizient:

Messwert Iceberg zu Kafka (nur Anhängen, genau einmal)
Worker-Maschinentyp e2-standard-4
Geschätzte vCPUs für 1 GB/s-Eingabe ca. 30 vCPUs
Geschätzte Worker für 1 GB/s-Eingabe ca. 7 Worker
Geschätzte SECU pro Stunde für 1 GB/s ca.0,2 SECU/h

Wichtige Erkenntnisse für umgekehrte Pipelines

  • Deutlich geringerer Rechenaufwand:Das Lesen und Projizieren von CDC-Streams aus Iceberg erfordert wesentlich weniger Rechenressourcen (ca. 30 vCPUs gegenüber ca. 450 vCPUs für direkte Schreibvorgänge), da die aufwendigen Vorgänge des Partitionierens, Codierens und Commit-Ausführens großer Mengen von Parquet-Dateien im Objektspeicher entfallen.
  • Ressourceneffizienz:Der nachgelagerte ereignisgesteuerte Verbrauch oder die Replikation aus einem Lakehouse-Format zurück in Streaming-Ebenen ist im Vergleich zu eingehenden Aufnahmepfaden sehr effizient.

Zusammenfassung der Architekturempfehlungen

Architekturmuster P99-Schreiblatenz Dateilayout Überlegungen zum Lesen nachgelagerter Daten
Kafka zu BigQuery (map_only) ca.5,4 s – Optimal (verwaltete BigQuery Storage Engine)
Kafka zu Iceberg mit der verwalteten BigQuery API ca.2,7 s Willkürlich kleine Dateien Regelmäßige Komprimierung erforderlich für Lesevorgänge mit hohem Volumen
Kafka zu Iceberg direkt (autosharding=false) ca.14,0 s ca. 800 KB Gut (größere anfängliche Dateigrößen, geringerer Komprimierungsbedarf)
Kafka zu Iceberg direkt (groupbykey) ca.13,0 s ca. 100 KB Mäßig (höherer Rechen- und Status-Overhead)

Kosten schätzen

Sie können die Referenzkosten Ihrer eigenen, vergleichbaren Pipeline mit ressourcenbasierter Abrechnung mit dem Google Cloud Preisrechner schätzen. Gehen Sie dazu so vor:

  1. Öffnen Sie den Preisrechner.
  2. Klicken Sie auf Der Schätzung hinzufügen.
  3. Wählen Sie Dataflow aus.
  4. Wählen Sie für Diensttyp die Option "Dataflow Classic" aus.
  5. Wählen Sie Erweiterte Einstellungen aus, um alle Optionen zu sehen.
  6. Wählen Sie den Standort aus, an dem der Job ausgeführt wird.
  7. Wählen Sie für Jobtyp die Option „Streaming“ aus.
  8. Wählen Sie Streaming Engine aktivieren aus.
  9. Geben Sie Informationen zu den Jobausführungsstunden, Worker-Knoten, Worker-Maschinen und zum nichtflüchtigen Speicher ein.
  10. Geben Sie die geschätzte Anzahl der Streaming Engine-Recheneinheiten ein.

Die Ressourcennutzung und die Kosten skalieren ungefähr linear mit dem Eingabedurchsatz. Bei kleinen Jobs mit nur wenigen Workern werden die Gesamtkosten jedoch von den Fixkosten dominiert. Als Ausgangspunkt können Sie die Anzahl der Worker-Knoten und den Ressourcenverbrauch aus den Benchmark-Ergebnissen ableiten.

Angenommen, Sie führen eine Pipeline mit der Architektur Kafka zu Iceberg direkt (autosharding=false) mit einer Eingabedatenrate von 100 MB/s aus. Anhand der Benchmark-Ergebnisse für eine 1-GB/s-Pipeline können Sie die Ressourcenanforderungen so schätzen:

  • Skalierungsfaktor: (100 MB/s) / (1.024 MB/s) = ca. 0,1
  • Geschätzte Worker-Knoten: 110 Worker × 0,1 = ca.11 Worker
  • Geschätzte Anzahl der Streaming Engine-Recheneinheiten pro Stunde: 200 × 0,1 = ca.20 Einheiten pro Stunde

Dieser Wert sollte nur als erste Schätzung verwendet werden. Der tatsächliche Durchsatz und die tatsächlichen Kosten können je nach Faktoren wie Maschinentyp, Verteilung der Nachrichtengröße, Nutzercode, Aggregationstyp, Schlüsselparallelität und Fenstergröße erheblich variieren. Weitere Informationen finden Sie unter Best Practices für die Kostenoptimierung von Dataflow.

Testpipeline ausführen

Verwenden Sie den gcloud dataflow flex-template run Befehl, um einen Apache Iceberg-Streamingjob mit der flexiblen Dataflow Vorlage bereitzustellen.

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'

Ersetzen Sie Folgendes:

  • JOB_NAME: der Name Ihres Dataflow-Jobs
  • PROJECT_ID: Ihre Google Cloud Projekt-ID
  • REGION: die Google Cloud Region, in der Ihr Job ausgeführt wird (z. B. us-central1)
  • KAFKA_BOOTSTRAP_ADDRESS: die Bootstrap-Adresse Ihres Apache Kafka-Clusters
  • KAFKA_TOPIC: der Name Ihres Kafka-Themas
  • ICEBERG_TABLE_IDENTIFIER: die Kennung Ihrer Ziel-Iceberg-Tabelle
  • CATALOG_NAME: der Name Ihres Iceberg-Katalogs
  • CATALOG_TYPE: der zu verwendende Katalogtyp (z. B. hadoop oder bigquery)
  • BUCKET_NAME: der Name des Cloud Storage-Bucket für Ihren Warehouse-Standort
  • SCHEMA_DEFINITION: die Schemadefinition für Ihre Kafka-Themendaten (z. B. {"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})