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.
- Kafka zu BigQuery (
map_only) (Referenz für die Leistung von Kafka zu BigQuery) - Kafka zu Iceberg direkt (
map_only,autosharding=false) - 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 Verteilungsmodushash.
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
- 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). - 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:
- Öffnen Sie den Preisrechner.
- Klicken Sie auf Der Schätzung hinzufügen.
- Wählen Sie Dataflow aus.
- Wählen Sie für Diensttyp die Option "Dataflow Classic" aus.
- Wählen Sie Erweiterte Einstellungen aus, um alle Optionen zu sehen.
- Wählen Sie den Standort aus, an dem der Job ausgeführt wird.
- Wählen Sie für Jobtyp die Option „Streaming“ aus.
- Wählen Sie Streaming Engine aktivieren aus.
- Geben Sie Informationen zu den Jobausführungsstunden, Worker-Knoten, Worker-Maschinen und zum nichtflüchtigen Speicher ein.
- 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-JobsPROJECT_ID: Ihre Google Cloud Projekt-IDREGION: die Google Cloud Region, in der Ihr Job ausgeführt wird (z. B.us-central1)KAFKA_BOOTSTRAP_ADDRESS: die Bootstrap-Adresse Ihres Apache Kafka-ClustersKAFKA_TOPIC: der Name Ihres Kafka-ThemasICEBERG_TABLE_IDENTIFIER: die Kennung Ihrer Ziel-Iceberg-TabelleCATALOG_NAME: der Name Ihres Iceberg-KatalogsCATALOG_TYPE: der zu verwendende Katalogtyp (z. B.hadoopoderbigquery)BUCKET_NAME: der Name des Cloud Storage-Bucket für Ihren Warehouse-StandortSCHEMA_DEFINITION: die Schemadefinition für Ihre Kafka-Themendaten (z. B.{"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})