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):
- Da Kafka a BigQuery (
map_only) (benchmark di base del rendimento da Kafka a BigQuery) - Da Kafka a Iceberg Direct (
map_only,autosharding=false) - 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 distribuzionehash.
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
- Overhead di Streaming Engine: l'aggiunta di un passaggio
groupbykeycon 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). - 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:
- Apri il calcolatore prezzi.
- Fai clic su Aggiungi alla stima.
- Seleziona Dataflow.
- Per Tipo di servizio, seleziona "Dataflow Classic".
- Seleziona Impostazioni avanzate per visualizzare l'insieme completo di opzioni.
- Scegli la località in cui viene eseguito il job.
- Per Tipo di prestazione, seleziona "Streaming".
- Seleziona Abilita Streaming Engine.
- Inserisci le informazioni relative alle ore di esecuzione del job, ai nodi worker, alle macchine worker e allo spazio di archiviazione su Persistent Disk.
- 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 DataflowPROJECT_ID: il tuo Google Cloud ID progettoREGION: la Google Cloud regione in cui viene eseguito il job (ad esempio,us-central1)KAFKA_BOOTSTRAP_ADDRESS: l'indirizzo di bootstrap del cluster Apache KafkaKAFKA_TOPIC: il nome dell'argomento KafkaICEBERG_TABLE_IDENTIFIER: l'identificatore della tabella Iceberg di destinazioneCATALOG_NAME: il nome del catalogo IcebergCATALOG_TYPE: il tipo di catalogo da utilizzare (ad esempio,hadoopobigquery)BUCKET_NAME: il nome del bucket Cloud Storage per la località del data warehouseSCHEMA_DEFINITION: la definizione dello schema per i dati dell'argomento Kafka (ad esempio,{"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})