Caratteristiche delle prestazioni delle pipeline da Kafka a BigQuery

Questa pagina descrive le caratteristiche delle prestazioni dei job di streaming Dataflow che leggono da Apache Kafka e scrivono in BigQuery. Fornisce i risultati dei test di benchmark per le pipeline di sola mappatura, che eseguono trasformazioni per messaggio senza tenere traccia dello stato o raggruppare gli elementi nello stream.

Molti carichi di lavoro di integrazione dei dati, inclusi ETL, convalida dei campi e mappatura degli schemi, rientrano nella categoria di sola mappatura. Se la tua pipeline segue questo pattern, puoi utilizzare questi benchmark per valutare il tuo job Dataflow rispetto a una configurazione di riferimento con buone prestazioni.

Metodologia di test

I benchmark sono stati condotti utilizzando le seguenti risorse:

  • Un cluster Managed Service per Apache Kafka. I messaggi sono stati generati utilizzando il modello Generatore di dati di streaming.

    • Frequenza dei messaggi: circa 1.000.000 di messaggi al secondo
    • Carico di input: 1 GiB/s
    • Formato del messaggio: testo JSON generato in modo casuale con uno schema fisso
    • Dimensioni del messaggio: circa 1 KiB per messaggio
    • Partizioni Kafka: 1000
  • Una tabella BigQuery standard.

  • Una pipeline di streaming Dataflow che utilizza il modello Apache Kafka a BigQuery. Questa pipeline esegue l'analisi e la mappatura degli schemi minime richieste. Non è stata utilizzata alcuna funzione definita dall'utente (UDF) personalizzata.

Dopo che la scalabilità orizzontale si è stabilizzata e la pipeline ha raggiunto lo stato stazionario, le pipeline sono state eseguite per circa un giorno, dopodiché i risultati sono stati raccolti e analizzati.

Pipeline Dataflow

Questo benchmark utilizza una pipeline di sola mappatura che esegue una semplice mappatura e conversione dei messaggi JSON. La pipeline è stata testata utilizzando sia la modalità "exactly-once" sia la modalità "at-least-once". L'elaborazione "at-least-once" offre una velocità effettiva migliore. Tuttavia, deve essere utilizzata solo quando i record duplicati sono accettabili o il sink downstream gestisce la deduplicazione.

Configurazione job

La tabella seguente mostra come sono stati configurati i job Dataflow.

Impostazione Valore
Tipo di macchina worker e2-standard-2
vCPU della macchina worker 2
RAM della macchina worker 8 GB
Persistent Disk della macchina worker Disco permanente standard (HDD), 30 GB
Numero massimo di worker 120
Streaming Engine
Scalabilità automatica orizzontale
Modello di fatturazione Fatturazione basata sulle risorse
L'API Storage Write (gRPC) è abilitata?
Stream dell'API Storage Write (gRPC) 400
Frequenza di attivazione dell'API Storage Write (gRPC) 5 secondi
Formato del messaggio JSON
Modalità di autenticazione Kafka

Credenziali predefinite dell'applicazione (ADC).

Per ulteriori informazioni, consulta Tipi di autenticazione per i broker Kafka.

L'API BigQuery Storage Write (gRPC) è consigliata per le pipeline di streaming. Quando utilizzi la modalità "exactly-once" con l'API Storage Write (gRPC), puoi modificare le seguenti impostazioni:

Per ulteriori informazioni, consulta Scrivere da Dataflow a BigQuery.

È necessario prestare particolare attenzione anche al numero di partizioni Apache Kafka. Per garantire un parallelismo delle chiavi sufficiente nella fase di lettura, il numero di partizioni deve essere almeno uguale al numero totale di vCPU worker. Per ulteriori informazioni, consulta Leggere da Apache Kafka a Dataflow.

Risultati benchmark

In questa sezione vengono descritti i risultati dei test di benchmark.

Velocità effettiva e utilizzo delle risorse

La tabella seguente mostra i risultati dei test per la velocità effettiva della pipeline e l'utilizzo delle risorse.

Risultato "Exactly-once" "At-least-once"
Velocità effettiva di input per worker Media: 15 MBps, n=3 Media: 18 MBps, n=3
Utilizzo medio della CPU in tutti i worker Media: 70%, n=3 Media: 75%, n=3
Numero di nodi worker Media: 63, n=3 Media: 53, n=3
Unità di calcolo di Streaming Engine all'ora Media: 58, n=3 Media: 0, n=3

L'algoritmo di scalabilità automatica può influire sul livello di utilizzo della CPU target. Per ottenere un utilizzo della CPU target più alto o più basso, puoi impostare l' intervallo di scalabilità automatica o il suggerimento per l'utilizzo dei worker. Gli obiettivi di utilizzo più elevati possono comportare costi inferiori, ma anche una latenza di coda peggiore, soprattutto per carichi variabili.

Latenza

La tabella seguente mostra i risultati dei benchmark per la latenza della pipeline in modalità "exactly-once", esclusa la fase di input.

Latenza end-to-end totale della fase, esclusa la fase di input "Exactly-once"
P50 Media: 1200 ms, n=3
P95 Media: 3000 ms, n=3
P99 Media: 5400 ms, n=3

I test hanno misurato la latenza end-to-end per fase (la job/streaming_engine/stage_end_to_end_latencies metrica) in tre esecuzioni di test a lunga esecuzione. Questa metrica misura il tempo trascorso da Streaming Engine in ogni fase della pipeline. Include tutti i passaggi interni della pipeline, ad esempio:

  • Shuffle e accodamento dei messaggi per l'elaborazione
  • Il tempo di elaborazione effettivo, ad esempio la conversione dei messaggi in oggetti riga
  • Scrittura dello stato persistente, nonché tempo trascorso in coda per scrivere lo stato persistente

A causa di una limitazione della metrica, la latenza della fase di input non viene segnalata. Pertanto, non è inclusa nel totale.

I benchmark mostrati qui rappresentano una baseline. La latenza è molto sensibile alla complessità della pipeline. Le UDF personalizzate, le trasformazioni aggiuntive e la logica di finestra complessa possono aumentare la latenza.

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 di Streaming Engine.

L'utilizzo delle risorse e il costo aumentano in modo approssimativamente lineare con la velocità effettiva 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 di sola mappatura in modalità "exactly-once", con una velocità di dati di input di 100 MiB/s. In base ai risultati dei benchmark per una pipeline da 1 GiB/s, puoi stimare i requisiti delle risorse come segue:

  • Fattore di scalabilità: (100 MiB/s) / (1 GiB/s) = 0,1
  • Nodi worker previsti: 63 worker × 0,1 = 6,3 worker
  • Numero previsto di unità di calcolo di Streaming Engine all'ora: 58 × 0,1 = 5,8 unità all'ora

Questo valore deve essere utilizzato solo come stima iniziale. La velocità effettiva 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

In questa sezione vengono mostrati i gcloud dataflow flex-template run comandi utilizzati per eseguire la pipeline di sola mappatura.

Modalità "exactly-once"

gcloud dataflow flex-template run JOB_NAME \
  --project=PROJECT_ID \
  --template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_to_BigQuery_Flex \
  --enable-streaming-engine \
  --parameters \
readBootstrapServerAndTopic="KAFKA_BOOTSTRAP_ADDRESS;KAFKA_TOPIC",\
kafkaReadAuthenticationMode=APPLICATION_DEFAULT_CREDENTIALS,\
messageFormat=JSON,\
writeMode=SINGLE_TABLE_NAME,\
outputTableSpec="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME",\
useBigQueryDLQ=true,\
outputDeadletterTable="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME_dlq",\
numStorageWriteApiStreams=400

Modalità "at-least-once"

gcloud dataflow flex-template run JOB_NAME \
  --project=PROJECT_ID \
  --template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_to_BigQuery_Flex \
  --enable-streaming-engine \
  --additional-experiments=streaming_mode_at_least_once \
  --parameters \
readBootstrapServerAndTopic="KAFKA_BOOTSTRAP_ADDRESS;KAFKA_TOPIC",\
kafkaReadAuthenticationMode=APPLICATION_DEFAULT_CREDENTIALS,\
messageFormat=JSON,\
writeMode=SINGLE_TABLE_NAME,\
outputTableSpec="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME",\
useBigQueryDLQ=true,\
outputDeadletterTable="PROJECT_ID:BQ_DATASET.BQ_TABLE_NAME_dlq",\
numStorageWriteApiStreams=400,\
useStorageWriteApiAtLeastOnce=true

Sostituisci quanto segue:

  • JOB_NAME: il nome del job Dataflow
  • PROJECT_ID: l'ID progetto
  • KAFKA_BOOTSTRAP_ADDRESS: l'indirizzo bootstrap del cluster Apache Kafka
  • KAFKA_TOPIC: il nome dell'argomento Kafka
  • BQ_DATASET: il nome del set di dati BigQuery
  • BQ_TABLE_NAME: il nome della tabella BigQuery

Genera dati di test

Per generare dati di test, utilizza il seguente comando per eseguire il modello Generatore di dati di streaming:

gcloud dataflow flex-template run JOB_NAME \
  --project=PROJECT_ID \
  --template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Streaming_Data_Generator \
  --max-workers=140 \
  --parameters \
schemaLocation=SCHEMA_LOCATION,\
qps=1000000,\
sinkType=KAFKA,\
bootstrapServer=KAFKA_BOOTSTRAP_ADDRESS,\
kafkaTopic=KAFKA_TOPIC,\
outputType=JSON

Sostituisci quanto segue:

  • JOB_NAME: il nome del job Dataflow
  • PROJECT_ID: l'ID progetto
  • SCHEMA_LOCATION: il percorso di un file di schema in Cloud Storage
  • KAFKA_BOOTSTRAP_ADDRESS: l'indirizzo bootstrap del cluster Apache Kafka
  • KAFKA_TOPIC: il nome dell'argomento Kafka

Il modello Generatore di dati di streaming utilizza un file JSON Data Generator per definire lo schema dei messaggi. I test di benchmark hanno utilizzato uno schema di messaggi simile al seguente:

{
  "logStreamId": "{{integer(1000001,2000000)}}",
  "message": "{{alphaNumeric(962)}}"
}

Passaggi successivi