Caratteristiche delle prestazioni delle pipeline da Kafka a Iceberg

Questa pagina descrive le caratteristiche di rendimento, a partire dalla versione 2.75.0 di Apache Beam, per i job di streaming di Dataflow che leggono da Apache Kafka e scrivono nelle tabelle Apache Iceberg. Valuta le differenze di rendimento tra le scritture dirette di Apache Iceberg e le scritture indirizzate tramite l'API BigQuery gestita e confronta questi risultati con i benchmark di base delle pipeline da Kafka a BigQuery. Poiché le ottimizzazioni per l'I/O di Apache Iceberg sono in corso, queste metriche di rendimento sono soggette a modifiche.

I confronti dei benchmark sono disponibili in tre configurazioni di mapping senza stato principali (il che significa che leggono dall'origine, convertono il messaggio in un record e scrivono nel sink senza tenere traccia dello stato o applicare una logica di business complessa; indicate come map_only o mapping nei benchmark):

  1. Da Kafka a BigQuery (map_only) (benchmark di base del rendimento da Kafka a BigQuery)
  2. Da Kafka a Iceberg Direct (map_only, autosharding=false)
  3. Da Kafka a Iceberg utilizzando l'API BigQuery gestita (map_only)

Inoltre, questa guida valuta i pattern di streaming diretti di Apache Iceberg, come il batch con stato utilizzando groupbykey, e descrive in dettaglio le considerazioni downstream cruciali relative alle distribuzioni delle dimensioni dei file, al comportamento di autosharding e alla latenza delle query lato lettura.

Metodologia di test

I benchmark sono stati condotti utilizzando le seguenti risorse:

  • Cluster Managed Service per Apache Kafka: il traffico è stato generato utilizzando il modello Dataflow Streaming Data Generator.
    • Throughput di input: 1 GBps
    • Frequenza dei messaggi: circa 1.000.000 di messaggi al secondo
    • Formato dei messaggi: testo JSON con uno schema fisso (circa 1 KB per messaggio)
    • Partizioni: 1000 partizioni Kafka
  • Sink di destinazione:
    • BigQuery: tabella standard (non partizionata) scritta utilizzando l'API BigQuery Storage Write.
    • Apache Iceberg: catalogo supportato da Cloud Storage. Il sink diretto è partizionato utilizzando bucket(id, 64) (suddiviso in 64 shard sulla chiave primaria) e utilizza la modalità di distribuzione hash.

Dopo la stabilizzazione della scalabilità automatica orizzontale, ogni configurazione della pipeline è stata eseguita in stato stazionario per 24 ore. I benchmark per ogni caso di pipeline sono stati eseguiti 3 volte separate e tutti i valori riportati rappresentano le medie calcolate in queste esecuzioni per garantire metriche di rendimento sostenute e affidabili.

Rendimento di importazione: carichi di lavoro di mapping

Le pipeline di mapping senza stato leggono dall'origine, convertono il formato del messaggio in un record e scrivono nel sink senza tenere traccia dello stato tra i record. Le sezioni seguenti analizzano le architetture di riferimento eseguite a 1 GBps.

Configurazione job

Impostazione Da Kafka a BigQuery (map_only) Da Kafka a Iceberg Direct (autosharding=false) Da Kafka a Iceberg utilizzando l'API BigQuery gestita
Tipo di macchina worker e2-standard-2 e2-standard-4 e2-standard-4
vCPU per worker 2 4 4
RAM per worker 8 GB 16 GB 16 GB
Streaming Engine Abilitato Abilitata Abilitato
Scalabilità automatica orizzontale Abilitato Abilitata Abilitato
Frequenza di attivazione 5 secondi 60 secondi 60 secondi

Throughput e utilizzo delle risorse

La scrittura diretta nei file Parquet fisici nell'archivio di oggetti comporta un overhead di I/O maggiore rispetto all'importazione di flussi di dati di BigQuery. Rispetto alle scritture dirette di Iceberg, l'indirizzamento delle scritture tramite l'API BigQuery gestita migliora l'utilizzo della CPU del worker (circa il 70% rispetto a circa il 60%) e riduce leggermente il consumo di Streaming Engine (circa 180 SECU/h rispetto a circa 200 SECU/h), anche se i requisiti di calcolo complessivi dei worker rimangono simili (circa 440 vCPU rispetto a circa 450 vCPU).

Metrica Da Kafka a BigQuery (map_only) Da Kafka a Iceberg Direct (autosharding=false) Da Kafka a Iceberg utilizzando l'API BigQuery gestita
Throughput di input medio per worker Circa 15 MBps Circa 9 MBps Circa 9 MBps
Utilizzo medio della CPU Circa 70% Circa 60% Circa 70%
vCPU stimate per input di 1 GBps Circa 126 vCPU Circa 450 vCPU Circa 440 vCPU
Worker stimati per input di 1 GBps Circa 63 worker Circa 110 worker Circa 110 worker
SECU stimate all'ora per 1 GBps Circa 58 SECU/h Circa 200 SECU/h Circa 180 SECU/h

Profilo della latenza di scrittura

Le scritture dirette di Iceberg mostrano una latenza di coda (P99) elevata a causa dei vincoli di commit dei metadati dell'archiviazione di oggetti. L'utilizzo dell'API BigQuery gestita elimina i picchi di latenza di coda mantenendo una latenza mediana bassa.

Latenza di scrittura end-to-end Da Kafka a BigQuery Da Kafka a Iceberg Direct (autosharding=false) Da Kafka a Iceberg utilizzando l'API BigQuery gestita
P50 (mediana) Circa 1200 ms Circa 1000 ms Circa 1000 ms
P95 Circa 3000 ms Circa 7400 ms Circa 1900 ms
P99 (coda) Circa 5400 ms Circa 14.000 ms Circa 2700 ms

Considerazioni e scelte di progettazione relative ad autosharding

Questa sezione illustra le implicazioni di autosharding sulle dimensioni dei file e sulla latenza della pipeline durante la scrittura in Apache Iceberg.

Perché autosharding=false è stato scelto come benchmark di base

Nei test iniziali, l'abilitazione di autosharding ha causato la riduzione delle dimensioni dei file in piccoli blocchi e la loro fluttuazione arbitraria a causa della suddivisione dinamica degli shard attivata da picchi di carico a livello di thread localizzati, anche con un carico di input aggregato costante.

Per mantenere layout di file Parquet stabili e prevedibili (circa 800 KB in media) e garantire un benchmark di base equo senza scaricamenti prematuri, è stato scelto autosharding=false per la configurazione del sink diretto.

Cosa succede se disabiliti autosharding rispetto a mantenerlo abilitato?

  • Con autosharding=false (benchmark di base): ottieni dimensioni dei file iniziali maggiori (circa 800 KB in media) rispetto ad autosharding. Sebbene sia ancora piccolo rispetto alle dimensioni ideali dei file Iceberg (128-512 MB), richiede una compattazione downstream significativamente inferiore. Tuttavia, il compromesso è una latenza di coda di scrittura elevata (P99 raggiunge circa 14,0 s) a causa dei colli di bottiglia dei metadati dell'archiviazione di oggetti.
  • Se autosharding è abilitato: Dataflow scala dinamicamente i thread di scrittura per assorbire i picchi di throughput locali, il che riduce la latenza di coda di scrittura. Tuttavia, compromette il livello di archiviazione producendo un volume elevato di file Parquet piccoli e frammentati (circa 100 KB o meno). Queste dimensioni dei file mostrano una varianza elevata e fluttuano arbitrariamente tra le esecuzioni (in media da circa 39 KB a circa 100 KB), aumentando la necessità di una manutenzione aggressiva della compattazione downstream.

Ottimizzazione e consigli per le partizioni

Durante la nostra valutazione, abbiamo sperimentato vari valori di partizione fissi per la tabella di destinazione per trovare un equilibrio ottimale. Abbiamo scoperto che l'utilizzo di 64 bucket (ad esempio, bucket(id, 64)) per il partizionamento della tabella di destinazione ha prodotto le dimensioni dei file target mantenendo un utilizzo e un throughput decenti. Questo approccio ci ha aiutato a ottenere i vantaggi in termini di rendimento di autosharding evitando i problemi di frammentazione arbitraria delle dimensioni dei file legati alla scalabilità completamente dinamica.

Consiglio per i professionisti: i clienti sono invitati a eseguire test preliminari simili con le impostazioni di partizione target per individuare il punto ottimale che massimizza il parallelismo della pipeline senza compromettere le dimensioni dei file Parquet.

Implicazioni di lettura downstream: dimensioni dei file e compattazione

Sebbene le metriche lato scrittura favoriscano l'API BigQuery gestita per l'importazione di Iceberg, l'efficienza complessiva della pipeline dipende in gran parte dal rendimento di lettura downstream:

  • Generazione di file di piccole dimensioni nell'API BigQuery gestita: l'API BigQuery gestita scarica i dati di frequente per garantire una latenza di scrittura bassa. Questo comportamento comporta un volume elevato di file Parquet di piccole dimensioni scritti nel catalogo Iceberg di destinazione.
  • Impatto sulla latenza delle query di lettura: i motori di query (ad esempio, Starburst/Trino, Apache Spark, BigQuery, Dremio) che leggono tabelle con milioni di file Parquet di piccole dimensioni comportano un overhead di analisi dei metadati elevato e penalità di scansione delle partizioni.
  • Requisiti di compattazione: per evitare il degrado del rendimento di lettura quando utilizzi l'API BigQuery gestita (o se autosharding è abilitato per le scritture dirette), esegui regolarmente i job di manutenzione della compattazione di Iceberg (ad esempio, REWRITE DATA FILES). L'overhead di calcolo per la compattazione deve essere incluso nella progettazione complessiva dell'architettura.
  • Distribuzione dei file di scrittura diretta (autosharding=false): le scritture dirette di Iceberg con sharding fisso producono file Parquet medi più grandi (circa 800 KB), con un layout meno frammentato per l'accesso immediato alle query senza richieste di compattazione immediate (anche se al di sotto dell'intervallo ideale).

Pipeline Iceberg dirette con stato (groupbykey)

Per valutare le strategie di batch manuale, il raggruppamento di chiavi con stato (groupbykey) è stato testato rispetto alla pipeline di base Da Kafka a Iceberg Direct (map_only, autosharding=false). Entrambe le configurazioni scrivono i file Parquet direttamente nell'archiviazione di oggetti.

Confronto con benchmark

Metrica / funzionalità Benchmark di base del sink diretto (autosharding=false) Sink diretto con stato (groupbykey) Impatto sul rendimento
vCPU stimate per 1 GBps Circa 450 vCPU Circa 520 vCPU Circa +16% di calcolo richiesto
Utilizzo medio della CPU Circa 60% Circa 50% Circa -17% di efficienza del worker
SECU/h stimate per 1 GBps Circa 200 SECU/h Circa 300 SECU/h Circa +50% di carico di Streaming Engine
Dimensioni medie dei file Circa 800 KB Circa 100 KB Genera batch di file più piccoli
Latenza P50 Circa 1000 ms Circa 1200 ms Circa +20% di mediana più lenta
Latenza P95 Circa 7400 ms Circa 5500 ms Circa -26% di latenza inferiore
Latenza P99 Circa 14.000 ms Circa 13.000 ms Variazione marginale della latenza di coda

Analisi compromesso

  1. Overhead di Streaming Engine: l'aggiunta di un passaggio groupbykey con stato richiede a Beam di archiviare lo stato intermedio tra i limiti delle finestre. Ciò aumenta il consumo di unità di calcolo Streaming Engine di circa il 50% (da circa 200 SECU/h a circa 300 SECU/h).
  2. Latenza di buffering: l'aggregazione manuale delle chiavi introduce il buffering delle finestre obbligatorio, aumentando la latenza di scrittura mediana (P50) a circa 1200 ms e la latenza P95 a circa 5,5 s.

Pipeline inverse: streaming da Iceberg a Kafka

Per valutare le funzionalità bidirezionali di lakehouse, sono stati condotti benchmark anche per i dati di streaming che scorrono in senso inverso: lettura di stream di sola aggiunta da una tabella Apache Iceberg e pubblicazione di nuovo in Apache Kafka.

Configurazione ed efficienza dei job

A differenza delle pipeline di importazione che devono gestire scritture di file di archiviazione di oggetti pesanti o colli di bottiglia di commit dei metadati, la lettura e lo streaming delle modifiche da Iceberg operano con elevata efficienza:

Metrica Da Iceberg a Kafka (di sola aggiunta, exactly-once)
Tipo di macchina worker e2-standard-4
vCPU stimate per input di 1 GBps Circa 30 vCPU
Worker stimati per input di 1 GBps Circa 7 worker
SECU stimate all'ora per 1 GBps Circa 0,2 SECU/h

Key takeaways per le pipeline inverse

  • Overhead di calcolo significativamente inferiore: la lettura e la proiezione degli stream CDC da Iceberg richiedono molte meno risorse di calcolo (circa 30 vCPU rispetto a circa 450 vCPU per le scritture dirette) perché evita il lavoro pesante di partizionamento, codifica e commit di grandi volumi di file Parquet nell'archiviazione di oggetti.
  • Efficienza delle risorse: il consumo o la replica downstream basati su eventi da un formato lakehouse ai livelli di streaming sono molto efficienti rispetto ai percorsi di importazione in entrata.

Riepilogo dei consigli sull'architettura

Pattern di architettura Latenza di scrittura P99 Layout dei file Considerazioni sulla lettura downstream
Da Kafka a BigQuery (map_only) Circa 5,4 s N/D Ottimale (motore di archiviazione BigQuery gestito)
Da Kafka a Iceberg utilizzando l'API BigQuery gestita Circa 2,7 s File arbitrariamente piccoli Richiede una compattazione periodica per le letture di volumi elevati
Da Kafka a Iceberg Direct (autosharding=false) Circa 14,0 s Circa 800 KB Buono (dimensioni dei file iniziali maggiori, minore richiesta di compattazione)
Da Kafka a Iceberg Direct (groupbykey) Circa 13,0 s Circa 100 KB Moderato (overhead di calcolo e stato più elevato)

Stima i costi

Puoi stimare il costo di base della tua pipeline comparabile con la fatturazione basata sulle risorse utilizzando il Google Cloud calcolatore prezzi, come segue:

  1. Apri il calcolatore prezzi.
  2. Fai clic su Aggiungi alla stima.
  3. Seleziona Dataflow.
  4. Per Tipo di servizio, seleziona "Dataflow Classic".
  5. Seleziona Impostazioni avanzate per visualizzare l'insieme completo di opzioni.
  6. Scegli la località in cui viene eseguito il job.
  7. Per Tipo di prestazione, seleziona "Streaming".
  8. Seleziona Abilita Streaming Engine.
  9. Inserisci le informazioni relative alle ore di esecuzione del job, ai nodi worker, alle macchine worker e allo spazio di archiviazione su Persistent Disk.
  10. Inserisci il numero stimato di unità di calcolo Streaming Engine.

L'utilizzo delle risorse e il costo aumentano in modo approssimativamente lineare con il throughput di input, anche se per i job di piccole dimensioni con solo pochi worker, il costo totale è dominato dai costi fissi. Come punto di partenza, puoi estrapolare il numero di nodi worker e il consumo di risorse dai risultati dei benchmark.

Supponiamo, ad esempio, di eseguire una pipeline utilizzando l'architettura Da Kafka a Iceberg Direct (autosharding=false) , con una velocità di dati di input di 100 MBps. In base ai risultati dei benchmark per una pipeline da 1 GBps, puoi stimare i requisiti delle risorse come segue:

  • Fattore di scalabilità: (100 MBps) / (1024 MBps) = circa 0,1
  • Nodi worker previsti: 110 worker × 0,1 = circa 11 worker
  • Numero previsto di unità di calcolo Streaming Engine all'ora: 200 × 0,1 = circa 20 unità all'ora

Questo valore deve essere utilizzato solo come stima iniziale. Il throughput e il costo effettivi possono variare in modo significativo in base a fattori quali il tipo di macchina, la distribuzione delle dimensioni dei messaggi, il codice utente, il tipo di aggregazione, il parallelismo delle chiavi e le dimensioni della finestra. Per ulteriori informazioni, consulta Best practice per l'ottimizzazione dei costi di Dataflow.

Esegui una pipeline di test

Per eseguire il deployment di un job di streaming di Apache Iceberg utilizzando il modello flessibile di Dataflow, utilizza il gcloud dataflow flex-template run comando.

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'

Sostituisci quanto segue:

  • JOB_NAME: il nome del job Dataflow
  • PROJECT_ID: il tuo Google Cloud ID progetto
  • REGION: la Google Cloud regione in cui viene eseguito il job (ad esempio, us-central1)
  • KAFKA_BOOTSTRAP_ADDRESS: l'indirizzo di bootstrap del cluster Apache Kafka
  • KAFKA_TOPIC: il nome dell'argomento Kafka
  • ICEBERG_TABLE_IDENTIFIER: l'identificatore della tabella Iceberg di destinazione
  • CATALOG_NAME: il nome del catalogo Iceberg
  • CATALOG_TYPE: il tipo di catalogo da utilizzare (ad esempio, hadoop o bigquery)
  • BUCKET_NAME: il nome del bucket Cloud Storage per la località del data warehouse
  • SCHEMA_DEFINITION: la definizione dello schema per i dati dell'argomento Kafka (ad esempio, {"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})