Scrivere i dati da Kafka a BigQuery utilizzando Dataflow

Questa pagina mostra come utilizzare Dataflow per leggere i dati da Google Cloud Managed Service per Apache Kafka e scrivere i record in una tabella BigQuery. Questo tutorial utilizza il modello Da Apache Kafka a BigQuery per creare il job Dataflow.

Panoramica

Apache Kafka è una piattaforma open source per lo streaming di eventi. Kafka è di uso comune nelle architetture distribuite per consentire la comunicazione tra componenti a basso accoppiamento. Puoi utilizzare Dataflow per leggere gli eventi da Kafka, elaborarli e scrivere i risultati in una tabella BigQuery per ulteriori analisi.

Managed Service per Apache Kafka è un Google Cloud servizio che ti aiuta a eseguire cluster Kafka sicuri e scalabili.

Lettura degli eventi Kafka in BigQuery
Architettura basata sugli eventi che utilizza Apache Kafka

Autorizzazioni obbligatorie

Il service account worker Dataflow deve disporre dei seguenti ruoli Identity and Access Management (IAM):

  • Managed Kafka Client (roles/managedkafka.client)
  • Editor dati BigQuery (roles/bigquery.dataEditor)

Per saperne di più, consulta Sicurezza e autorizzazioni di Dataflow.

Crea un cluster Kafka

In questo passaggio creerai un cluster Managed Service per Apache Kafka. Per saperne di più, consulta Creare un cluster Managed Service per Apache Kafka.

Console

  1. Vai alla pagina Managed Service per Apache Kafka > Cluster.

    Vai a Cluster

  2. Fai clic su Crea.

  3. Nella casella Nome cluster, inserisci un nome per il cluster.

  4. Nell'elenco Regione, seleziona una località per il cluster.

  5. Fai clic su Crea.

gcloud

Utilizza il managed-kafka clusters create comando.

gcloud managed-kafka clusters create CLUSTER \
--location=REGION \
--cpu=3 \
--memory=3GiB \
--subnets=projects/PROJECT_ID/regions/REGION/subnetworks/SUBNET_NAME

Sostituisci quanto segue:

  • CLUSTER: un nome per il cluster
  • REGION: la regione in cui hai creato la subnet
  • PROJECT_ID: il tuo ID progetto
  • SUBNET_NAME: la subnet in cui vuoi eseguire il deployment del cluster

In genere, la creazione di un cluster richiede 20-30 minuti.

Crea un argomento Kafka

Dopo aver creato il cluster Managed Service per Apache Kafka, crea un argomento.

Console

  1. Vai alla pagina Managed Service per Apache Kafka > Cluster.

    Vai a Cluster

  2. Fai clic sul nome del cluster.

  3. Nella pagina dei dettagli del cluster, fai clic su Crea argomento.

  4. Nella casella Nome argomento, inserisci un nome per l'argomento.

  5. Fai clic su Crea.

gcloud

Utilizza il managed-kafka topics create comando.

gcloud managed-kafka topics create TOPIC_NAME \
--cluster=CLUSTER \
--location=REGION \
--partitions=10 \
--replication-factor=3

Sostituisci quanto segue:

  • TOPIC_NAME: il nome dell'argomento da creare

Crea una tabella BigQuery

In questo passaggio creerai una tabella BigQuery con il seguente schema:

Nome colonna Tipo di dati
name STRING
customer_id INTEGER

Se non hai già un set di dati BigQuery, creane uno. Per saperne di più, consulta Creare set di dati. Poi crea una nuova tabella vuota:

Console

  1. Vai alla pagina BigQuery.

    Vai a BigQuery

  2. Nel riquadro Explorer , espandi il progetto e seleziona un set di dati.

  3. Nella sezione delle informazioni del set di dati, fai clic su Crea tabella.

  4. Nell'elenco Crea tabella da, seleziona Tabella vuota.

  5. Nella casella Tabella, inserisci il nome della tabella.

  6. Nella sezione Schema, fai clic su Modifica come testo.

  7. Incolla la seguente definizione di schema:

    name:STRING,
    customer_id:INTEGER
    
  8. Fai clic su Crea tabella.

gcloud

Utilizza il bq mk comando.

bq mk --table \
  PROJECT_ID:DATASET_NAME.TABLE_NAME \
  name:STRING,customer_id:INTEGER

Sostituisci quanto segue:

  • PROJECT_ID: il tuo ID progetto
  • DATASET_NAME: il nome del set di dati
  • TABLE_NAME: il nome della tabella da creare

Esegui il job Dataflow

Dopo aver creato il cluster Kafka e la tabella BigQuery, esegui il modello Dataflow.

Console

Innanzitutto, recupera l'indirizzo del server bootstrap del cluster:

  1. Nella Google Cloud console, vai alla pagina Cluster.

    Vai a Cluster

  2. Fai clic sul nome del cluster.

  3. Fai clic sulla scheda Configurazioni.

  4. Copia l'indirizzo del server bootstrap da URL bootstrap.

Poi, esegui il modello per creare il job Dataflow:

  1. Vai alla pagina Dataflow > Job.

    Vai a Job

  2. Fai clic su Crea job da modello.

  3. Nel campo Nome job, inserisci kafka-to-bq.

  4. Per Endpoint regionale, seleziona la regione in cui si trova il tuo cluster Managed Service per Apache Kafka.

  5. Seleziona il modello "Da Kafka a BigQuery".

  6. Inserisci i seguenti parametri del modello:

    • Server bootstrap Kafka: l'indirizzo del server bootstrap
    • Argomento Kafka di origine: il nome dell'argomento da leggere
    • Modalità di autenticazione dell'origine Kafka: APPLICATION_DEFAULT_CREDENTIALS
    • Formato messaggi Kafka: JSON
    • Strategia per il nome della tabella: SINGLE_TABLE_NAME
    • Tabella di output BigQuery: la tabella BigQuery, formattata come segue: PROJECT_ID:DATASET_NAME.TABLE_NAME
  7. In Coda di messaggi non recapitabili, seleziona Scrivi errori in BigQuery.

  8. Inserisci un nome di tabella BigQuery per la coda di messaggi non recapitabili, formattato come segue: PROJECT_ID:DATASET_NAME.ERROR_TABLE_NAME

    Non creare questa tabella in anticipo. La pipeline la crea.

  9. Fai clic su Esegui job.

gcloud

Utilizza il dataflow flex-template run comando.

gcloud dataflow flex-template run kafka-to-bq \
--template-file-gcs-location gs://dataflow-templates/latest/flex/Kafka_to_BigQuery \
--region LOCATION \
--parameters \
readBootstrapServerAndTopic=projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID/topics/TOPIC,\
persistKafkaKey=false,\
writeMode=SINGLE_TABLE_NAME,\
kafkaReadOffset=earliest,\
kafkaReadAuthenticationMode=APPLICATION_DEFAULT_CREDENTIALS,\
messageFormat=JSON,\
outputTableSpec=PROJECT_ID:DATASET_NAME.TABLE_NAME\
useBigQueryDLQ=true,\
outputDeadletterTable=PROJECT_ID:DATASET_NAME.ERROR_TABLE_NAME

Sostituisci le seguenti variabili:

  • LOCATION: la regione in cui si trova Managed Service per Apache Kafka
  • PROJECT_ID: il nome del Google Cloud progetto
  • CLUSTER_ID: il nome del cluster
  • TOPIC: il nome dell'argomento Kafka
  • DATASET_NAME: il nome del set di dati
  • TABLE_NAME: il nome della tabella
  • ERROR_TABLE_NAME: un nome di tabella BigQuery per la coda di messaggi non recapitabili

Non creare la tabella per la coda di messaggi non recapitabili in anticipo. La pipeline la crea.

Invia messaggi a Kafka

Dopo l'avvio del job Dataflow, puoi inviare messaggi a Kafka e la pipeline li scrive in BigQuery.

  1. Crea una VM nella stessa subnet del cluster Kafka e installa gli strumenti a riga di comando di Kafka. Per istruzioni dettagliate, consulta Configurare una macchina client in Pubblicare e utilizzare i messaggi con l'interfaccia a riga di comando.

  2. Esegui il comando seguente per scrivere i messaggi nell'argomento Kafka:

    kafka-console-producer.sh \
     --topic TOPIC \
     --bootstrap-server bootstrap.CLUSTER_ID.LOCATION.managedkafka.PROJECT_ID.cloud.goog:9092 \
     --producer.config client.properties

    Sostituisci le seguenti variabili:

    • TOPIC: il nome dell'argomento Kafka
    • CLUSTER_ID: il nome del cluster
    • LOCATION: la regione in cui si trova il cluster
    • PROJECT_ID: il nome del Google Cloud progetto
  3. Al prompt, inserisci le seguenti righe di testo per inviare messaggi a Kafka:

    {"name": "Alice", "customer_id": 1}
    {"name": "Bob", "customer_id": 2}
    {"name": "Charles", "customer_id": 3}
    

Utilizza una coda di messaggi non recapitabili

Durante l'esecuzione del job, la pipeline potrebbe non riuscire a scrivere singoli messaggi in BigQuery. I possibili errori includono:

  • Errori di serializzazione, incluso JSON in formato errato.
  • Errori di conversione del tipo, causati da una mancata corrispondenza tra lo schema della tabella e i dati JSON.
  • Campi aggiuntivi nei dati JSON che non sono presenti nello schema della tabella.

Questi errori non causano l'errore del job e non vengono visualizzati come errori nel log del job Dataflow. La pipeline utilizza invece una coda di messaggi non recapitabili per gestire questi tipi di errori.

Per abilitare la coda di messaggi non recapitabili quando esegui il modello, imposta i seguenti parametri del modello:

  • useBigQueryDLQ: true
  • outputDeadletterTable: un nome di tabella BigQuery completo; ad esempio, my-project:dataset1.errors

La pipeline crea automaticamente la tabella. Se si verifica un errore durante l'elaborazione di un messaggio Kafka, la pipeline scrive una voce di errore nella tabella.

Esempi di messaggi di errore:

Tipo di errore Dati sull'evento errorMessage
Errore di serializzazione "Hello world" Failed to serialize json to table row: "Hello world"
Errore di conversione del tipo {"name":"Emily","customer_id":"abc"} { "errors" : [ { "debugInfo" : "", "location" : "age", "message" : "Cannot convert value to integer (bad value): abc", "reason" : "invalid" } ], "index" : 0 }
Campo sconosciuto {"name":"Zoe","age":34} { "errors" : [ { "debugInfo" : "", "location" : "age", "message" : "no such field: customer_id.", "reason" : "invalid" } ], "index" : 0 }

Utilizza i tipi di dati BigQuery

Internamente, il connettore I/O Kafka converte i payload dei messaggi JSON in oggetti TableRow di Apache Beam e traduce i valori dei campi TableRow in tipi BigQuery.

La tabella seguente mostra le rappresentazioni JSON dei tipi di dati BigQuery .

Tipo BigQuery Rappresentazione JSON
ARRAY [1.2,3]
BOOL true
DATE "2022-07-01"
DATETIME "2022-07-01 12:00:00.00"
DECIMAL 5.2E11
FLOAT64 3.142
GEOGRAPHY "POINT(1 2)"

Specifica la geografia utilizzando il formato WKT (Well-Known Text) o GeoJSON, formattato come stringa. Per saperne di più, consulta Caricare dati geospaziali.

INT64 10
INTERVAL "0-13 370 48:61:61"
STRING "string_val"
TIMESTAMP "2022-07-01T12:00:00.00Z"

Utilizza il metodo Date.toJSON di JavaScript per formattare il valore.

Dati strutturati

Se i messaggi JSON seguono uno schema coerente, puoi rappresentare gli oggetti JSON utilizzando il STRUCT tipo di dati in BigQuery.

Nell'esempio seguente, il campo answers è un oggetto JSON con due sottocampi, a e b:

{"name":"Emily","answers":{"a":"yes","b":"no"}}

La seguente istruzione SQL crea una tabella BigQuery con uno schema compatibile:

CREATE TABLE my_dataset.kafka_events (name STRING, answers STRUCT<a STRING, b STRING>);

La tabella risultante sarà simile alla seguente:

+-------+----------------------+
| name  |       answers        |
+-------+----------------------+
| Emily | {"a":"yes","b":"no"} |
+-------+----------------------+

e semistrutturati

Se i messaggi JSON non seguono uno schema rigido, valuta la possibilità di archiviarli in BigQuery come tipo di dati JSON. Se archivi i dati JSON come tipo JSON, non devi definire lo schema in anticipo. Dopo importazione dati, puoi eseguire query sui dati utilizzando gli operatori di accesso ai campi (notazione con punti) e di accesso agli array in GoogleSQL. Per saperne di più, consulta Utilizzare i dati JSON in GoogleSQL.

Utilizza una funzione definita dall'utente per trasformare i dati

Questo tutorial presuppone che i messaggi Kafka siano formattati come JSON e che lo schema della tabella BigQuery corrisponda ai dati JSON, senza che siano state applicate trasformazioni ai dati.

Facoltativamente, puoi fornire una funzione definita dall'utente (UDF) JavaScript che trasforma i dati prima che vengano scritti in BigQuery. La funzione definita dall'utente può anche eseguire un'elaborazione aggiuntiva, come il filtraggio, la rimozione delle informazioni che consentono l'identificazione personale (PII) o l'arricchimento dei dati con campi aggiuntivi.

Per saperne di più, consulta Creare funzioni definite dall'utente per i modelli Dataflow.

Passaggi successivi