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.
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
Vai alla pagina Managed Service per Apache Kafka > Cluster.
Fai clic su Crea.
Nella casella Nome cluster, inserisci un nome per il cluster.
Nell'elenco Regione, seleziona una località per il cluster.
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 clusterREGION: la regione in cui hai creato la subnetPROJECT_ID: il tuo ID progettoSUBNET_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
Vai alla pagina Managed Service per Apache Kafka > Cluster.
Fai clic sul nome del cluster.
Nella pagina dei dettagli del cluster, fai clic su Crea argomento.
Nella casella Nome argomento, inserisci un nome per l'argomento.
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
Vai alla pagina BigQuery.
Nel riquadro Explorer , espandi il progetto e seleziona un set di dati.
Nella sezione delle informazioni del set di dati, fai clic su Crea tabella.
Nell'elenco Crea tabella da, seleziona Tabella vuota.
Nella casella Tabella, inserisci il nome della tabella.
Nella sezione Schema, fai clic su Modifica come testo.
Incolla la seguente definizione di schema:
name:STRING, customer_id:INTEGERFai 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 progettoDATASET_NAME: il nome del set di datiTABLE_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:
Nella Google Cloud console, vai alla pagina Cluster.
Fai clic sul nome del cluster.
Fai clic sulla scheda Configurazioni.
Copia l'indirizzo del server bootstrap da URL bootstrap.
Poi, esegui il modello per creare il job Dataflow:
Vai alla pagina Dataflow > Job.
Fai clic su Crea job da modello.
Nel campo Nome job, inserisci
kafka-to-bq.Per Endpoint regionale, seleziona la regione in cui si trova il tuo cluster Managed Service per Apache Kafka.
Seleziona il modello "Da Kafka a BigQuery".
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
In Coda di messaggi non recapitabili, seleziona Scrivi errori in BigQuery.
Inserisci un nome di tabella BigQuery per la coda di messaggi non recapitabili, formattato come segue:
PROJECT_ID:DATASET_NAME.ERROR_TABLE_NAMENon creare questa tabella in anticipo. La pipeline la crea.
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 KafkaPROJECT_ID: il nome del Google Cloud progettoCLUSTER_ID: il nome del clusterTOPIC: il nome dell'argomento KafkaDATASET_NAME: il nome del set di datiTABLE_NAME: il nome della tabellaERROR_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.
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.
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 KafkaCLUSTER_ID: il nome del clusterLOCATION: la regione in cui si trova il clusterPROJECT_ID: il nome del Google Cloud progetto
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:trueoutputDeadletterTable: 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 |
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.