Elaborare un flusso di modifiche Bigtable

Questo tutorial mostra come eseguire il deployment di una pipeline di dati su Dataflow per un flusso in tempo reale di modifiche del database provenienti dal flusso di modifiche di una tabella Bigtable. L'output della pipeline viene scritto in una serie di file su Cloud Storage.

Viene fornito un set di dati di esempio per un'applicazione di ascolto di musica. In questo tutorial, monitori i brani ascoltati e poi classifichi i primi cinque in un determinato periodo.

Questo tutorial è destinato agli utenti tecnici che hanno familiarità con la scrittura di codice e il deployment di pipeline di dati su Google Cloud.

Obiettivi

Questo tutorial mostra gli aspetti seguenti:

  • Crea una tabella Bigtable con un flusso di modifiche abilitato.
  • Esegui il deployment di una pipeline su Dataflow che trasforma e restituisce l'output del flusso di modifiche.
  • Visualizza i risultati della pipeline di dati.

Costi

In questo documento vengono utilizzati i seguenti componenti fatturabili di Google Cloud:

Per generare una stima dei costi in base all'utilizzo previsto, utilizza il calcolatore prezzi.

I nuovi Google Cloud utenti potrebbero avere diritto a una prova senza costi.

Al termine delle attività descritte in questo documento, puoi evitare l'addebito di ulteriori costi eliminando le risorse che hai creato. Per saperne di più, consulta Esegui la pulizia.

Prima di iniziare

    Installa Google Cloud CLI, quindi accedi a gcloud CLI con la tua identità federata. Dopo aver eseguito l'accesso, inizializza Google Cloud CLI eseguendo il comando seguente:

    gcloud init

    Crea o seleziona un Google Cloud progetto.

    Ruoli richiesti per selezionare o creare un progetto

    • Seleziona un progetto: la selezione di un progetto non richiede un ruolo IAM specifico: puoi selezionare qualsiasi progetto su cui ti è stato concesso un ruolo.
    • Crea un progetto: per creare un progetto, devi disporre del ruolo Autore progetto (roles/resourcemanager.projectCreator), che contiene l' resourcemanager.projects.create autorizzazione. Scopri come concedere i ruoli.
    • Crea un Google Cloud progetto:

      gcloud projects create PROJECT_ID

      Sostituisci PROJECT_ID con un nome per il Google Cloud progetto che stai creando.

    • Seleziona il Google Cloud progetto che hai creato:

      gcloud config set project PROJECT_ID

      Sostituisci PROJECT_ID con il nome del Google Cloud progetto.

    Verifica che la fatturazione sia attivata per il tuo Google Cloud progetto.

    Abilita le API Dataflow, API Cloud Bigtable, API Cloud Bigtable Admin e Cloud Storage:

    Ruoli richiesti per abilitare le API

    Per abilitare le API, devi disporre dell'autorizzazione serviceusage.services.enable. Se hai creato il progetto, probabilmente hai già questa autorizzazione tramite il ruolo Proprietario (roles/owner). In caso contrario, puoi ottenere questa autorizzazione tramite il ruolo Amministratore utilizzo servizi (roles/serviceusage.serviceUsageAdmin). Scopri come concedere i ruoli.

    gcloud services enable dataflow.googleapis.com bigtable.googleapis.com bigtableadmin.googleapis.com storage.googleapis.com
  1. Aggiorna e installa la cbt CLI .
    gcloud components update
    gcloud components install cbt

Prepara l'ambiente

Ottieni il codice

Clona il repository che contiene il codice campione. Se hai già scaricato questo repository in precedenza, esegui il pull per ottenere l'ultima versione.

git clone https://github.com/GoogleCloudPlatform/java-docs-samples.git
cd java-docs-samples/bigtable/beam/change-streams

Crea un bucket

  • Crea un bucket Cloud Storage:
    gcloud storage buckets create gs://BUCKET_NAME
    Sostituisci BUCKET_NAME con un nome del bucket che soddisfi i requisiti di denominazione dei bucket.
  • Crea un'istanza Bigtable

    Puoi utilizzare un'istanza esistente per questo tutorial o creare un'istanza con le configurazioni predefinite in una regione vicina a te.

    Crea una tabella

    L'applicazione di esempio monitora i brani ascoltati dagli utenti e archivia gli eventi di ascolto in Bigtable. Crea una tabella con un flusso di modifiche abilitato che abbia una famiglia di colonne (cf) e una colonna (song) e utilizzi gli ID utente per le chiavi di riga.

    Crea la tabella.

    gcloud bigtable instances tables create song-rank \
    --column-families=cf --change-stream-retention-period=7d \
    --instance=BIGTABLE_INSTANCE_ID --project=PROJECT_ID
    

    Sostituisci quanto segue:

    • PROJECT_ID: l'ID del progetto che stai utilizzando
    • BIGTABLE_INSTANCE_ID: l'ID dell'istanza che conterrà la nuova tabella

    Avvia la pipeline

    Questa pipeline trasforma il flusso di modifiche nel seguente modo:

    1. Legge il flusso di modifiche
    2. Recupera il titolo del brano
    3. Raggruppa gli eventi di ascolto dei brani in finestre di N secondi
    4. Conta i primi cinque brani
    5. Restituisce l'output dei risultati

    Esegui la pipeline.

    mvn compile exec:java -Dexec.mainClass=SongRank \
    "-Dexec.args=--project=PROJECT_ID --bigtableProjectId=PROJECT_ID \
    --bigtableInstanceId=BIGTABLE_INSTANCE_ID --bigtableTableId=song-rank \
    --outputLocation=gs://BUCKET_NAME/ \
    --runner=dataflow --region=BIGTABLE_REGION --experiments=use_runner_v2"
    

    Sostituisci BIGTABLE_REGION con l'ID della regione in cui si trova l'istanza Bigtable, ad esempio us-east5.

    Informazioni sulla pipeline

    I seguenti snippet di codice della pipeline possono aiutarti a comprendere il codice che stai eseguendo.

    Lettura del flusso di modifiche

    Il codice in questo esempio configura il flusso di origine con i parametri per la tabella e l'istanza Bigtable specifiche.

    p.apply(
            "Stream from Bigtable",
            BigtableIO.readChangeStream()
                .withProjectId(options.getBigtableProjectId())
                .withInstanceId(options.getBigtableInstanceId())
                .withTableId(options.getBigtableTableId())
                .withAppProfileId(options.getBigtableAppProfile())
    
        )

    Recupero del titolo del brano

    Quando viene ascoltato un brano, il titolo viene scritto nella famiglia di colonne cf e nel qualificatore di colonna song, quindi il codice estrae il valore dalla mutazione del flusso di modifiche e lo restituisce al passaggio successivo della pipeline.

    private static class ExtractSongName extends DoFn<KV<ByteString, ChangeStreamMutation>, String> {
    
      @DoFn.ProcessElement
      public void processElement(ProcessContext c) {
    
        for (Entry e : Objects.requireNonNull(Objects.requireNonNull(c.element()).getValue())
            .getEntries()) {
          if (e instanceof SetCell) {
            SetCell setCell = (SetCell) e;
            if ("cf".equals(setCell.getFamilyName())
                && "song".equals(setCell.getQualifier().toStringUtf8())) {
              c.output(setCell.getValue().toStringUtf8());
            }
          }
        }
      }
    }

    Conteggio dei primi cinque brani

    Puoi utilizzare le funzioni Beam integrate Count e Top.of per ottenere i primi cinque brani nella finestra corrente.

    .apply(Count.perElement())
    .apply("Top songs", Top.of(5, new SongComparator()).withoutDefaults())

    Output dei risultati

    Questa pipeline scrive i risultati nell'output standard e nei file. Per i file, le scritture vengono suddivise in gruppi di 10 elementi o segmenti di un minuto.

    .apply("Print", ParDo.of(new PrintFn()))
    .apply(
        "Collect at least 10 elements or 1 minute of elements",
        Window.<String>into(new GlobalWindows())
            .triggering(
                Repeatedly.forever(
                    AfterFirst.of(
                        AfterPane.elementCountAtLeast(10),
                        AfterProcessingTime
                            .pastFirstElementInPane()
                            .plusDelayOf(Duration.standardMinutes(1)
                            )
                    )
                ))
            .discardingFiredPanes())
    .apply(
        "Output top songs",
        TextIO.write()
            .to(options.getOutputLocation() + "song-charts/")
            .withSuffix(".txt")
            .withNumShards(1)
            .withWindowedWrites()
    );

    Visualizza la pipeline

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

      Vai a Dataflow

    2. Fai clic sul job con un nome che inizia con song-rank.

    3. Nella parte inferiore dello schermo, fai clic su Mostra per aprire il riquadro dei log.

    4. Fai clic su Log dei worker per monitorare i log di output del flusso di modifiche.

    Scritture di flussi

    Utilizza la CLI cbt per scrivere un numero di ascolti di brani per vari utenti nella tabella song-rank. Questa operazione è progettata per scrivere in pochi minuti per simulare lo streaming degli ascolti dei brani nel tempo.

    cbt -instance=BIGTABLE_INSTANCE_ID -project=PROJECT_ID import \
    song-rank song-rank-data.csv  column-family=cf batch-size=1
    

    Visualizza l'output

    Leggi l'output su Cloud Storage per visualizzare i brani più ascoltati.

    gcloud storage cat gs://BUCKET_NAME/song-charts/GlobalWindow-pane-0-00000-of-00001.txt
    

    Output di esempio:

    2023-07-06T19:53:38.232Z [KV{The Wheels on the Bus, 199}, KV{Twinkle, Twinkle, Little Star, 199}, KV{Ode to Joy , 192}, KV{Row, Row, Row Your Boat, 186}, KV{Take Me Out to the Ball Game, 182}]
    2023-07-06T19:53:49.536Z [KV{Old MacDonald Had a Farm, 20}, KV{Take Me Out to the Ball Game, 18}, KV{Für Elise, 17}, KV{Ode to Joy , 15}, KV{Mary Had a Little Lamb, 12}]
    2023-07-06T19:53:50.425Z [KV{Twinkle, Twinkle, Little Star, 20}, KV{The Wheels on the Bus, 17}, KV{Row, Row, Row Your Boat, 13}, KV{Happy Birthday to You, 12}, KV{Over the Rainbow, 9}]
    

    Libera spazio

    Per evitare che al tuo account Google Cloud vengano addebitati costi relativi alle risorse utilizzate in questo tutorial, elimina il progetto che contiene le risorse oppure mantieni il progetto ed elimina le singole risorse.

    Elimina il progetto

      Elimina un Google Cloud progetto:

      gcloud projects delete PROJECT_ID

    Elimina singole risorse

    1. Elimina il bucket e i file.

      gcloud storage rm --recursive gs://BUCKET_NAME/
      
    2. Disabilita il flusso di modifiche nella tabella.

      gcloud bigtable instances tables update song-rank --instance=BIGTABLE_INSTANCE_ID \
      --clear-change-stream-retention-period
      
    3. Elimina la tabella song-rank.

      cbt -instance=BIGTABLE_INSTANCE_ID -project=PROJECT_ID deletetable song-rank
      
    4. Arresta la pipeline del flusso di modifiche.

      1. Elenca i job per ottenere l'ID job.

        gcloud dataflow jobs list --region=BIGTABLE_REGION
        
      2. Annulla il job.

        gcloud dataflow jobs cancel JOB_ID --region=BIGTABLE_REGION
        

        Sostituisci JOB_ID con l'ID job visualizzato dopo il comando precedente.

    Passaggi successivi