Tutorial: implementare il modello di acquisizione del data lake per aggiornare una tabella Delta di Databricks

Questa esercitazione descrive come gestire eventi in un account di archiviazione che include uno spazio dei nomi gerarchico.

Costruisci una piccola soluzione che ti permette di popolare una tabella Delta di Databricks caricando un file di valori separati da virgole (CSV) che descrive un ordine di vendita. Si costruisce questa soluzione collegando un abbonamento a Event Grid, una funzione Azure e un Job in Azure Databricks.

In questa esercitazione, imparerai a:

  • Creare una sottoscrizione di Event Grid che richiama una funzione di Azure.
  • Creare una funzione di Azure che riceve una notifica da un evento e quindi esegue il processo in Azure Databricks.
  • Creare un processo Databricks che inserisce un ordine di un cliente in una tabella Databricks Delta situata nell'account di archiviazione.

Costruisci questa soluzione in ordine inverso, partendo dallo spazio di lavoro di Azure Databricks.

Prerequisiti

Creazione di un ordine cliente

Per prima cosa, crea un file CSV che descriva un ordine di vendita, e poi carica quel file sull'account di archiviazione. Successivamente, utilizzi i dati di questo file per popolare la prima riga della tua tabella Delta Databricks.

  1. Passare al nuovo account di archiviazione nel portale di Azure.

  2. Seleziona Storage browser>Contenitori BLOB>Aggiungi contenitore e crea un nuovo contenitore denominato data.

    Screenshot della creazione di un container nel browser Archiviazione di Azure.

  3. Nel contenitore di dati creare una directory denominata input.

  4. Incollare il testo seguente in un editor di testo.

    InvoiceNo,StockCode,Description,Quantity,InvoiceDate,UnitPrice,CustomerID,Country
    536365,85123A,WHITE HANGING HEART T-LIGHT HOLDER,6,12/1/2010 8:26,2.55,17850,United Kingdom
    
  5. Salva questo file sul tuo computer locale e chiamalo data.csv.

  6. Nel browser di archiviazione, carica questo file nella cartella di input .

Creare un processo in Azure Databricks

In questa sezione, esegui questi compiti:

  • Creare un'area di lavoro di Azure Databricks.
  • Crea un notebook.
  • Creare e popolare una tabella di Databricks Delta.
  • Aggiungere codice per l'inserimento di righe nella tabella di Databricks Delta.
  • Creare un processo.

Creare un'area di lavoro di Azure Databricks

In questa sezione crei uno spazio di lavoro Azure Databricks utilizzando il portale Azure.

  1. Creare un'area di lavoro di Azure Databricks. Nomina lo spazio di lavoro contoso-orders. Vedere Creare un'area di lavoro di Azure Databricks.

  2. Creare un cluster. Assegna al cluster il nome customer-order-cluster. Vedere Creare un cluster.

  3. Crea un notebook. Assegnare un nome al notebook configure-customer-table e scegliere Python come linguaggio predefinito del notebook. Vedere Creare un notebook.

Creare e popolare una tabella di Databricks Delta

  1. Nel notebook creato copiare e incollare il blocco di codice seguente nella prima cella, ma non eseguire ancora il codice.

    Sostituisci i valori segnaposto appId, password e tenant in questo blocco di codice con i valori che hai raccolto durante il completamento dei prerequisiti di questa esercitazione.

    dbutils.widgets.text('source_file', "", "Source File")
    
    spark.conf.set("fs.azure.account.auth.type", "OAuth")
    spark.conf.set("fs.azure.account.oauth.provider.type", "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider")
    spark.conf.set("fs.azure.account.oauth2.client.id", "<appId>")
    spark.conf.set("fs.azure.account.oauth2.client.secret", "<password>")
    spark.conf.set("fs.azure.account.oauth2.client.endpoint", "https://login.microsoftonline.com/<tenant>/oauth2/token")
    
    adlsPath = 'abfss://data@contosoorders.dfs.core.windows.net/'
    inputPath = adlsPath + dbutils.widgets.get('source_file')
    customerTablePath = adlsPath + 'delta-tables/customers'
    

    Questo codice crea un widget denominato source_file. Più avanti si creerà una funzione di Azure che chiama questo codice e passa un percorso di file al widget. Questo codice autentica anche l'entità servizio con l'account di archiviazione e crea alcune variabili da usare in altre celle.

    Nota

    In un ambiente di produzione è consigliabile archiviare la chiave di autenticazione in Azure Databricks. Poi, aggiungi una chiave di ricerca al blocco di codice invece della chiave di autenticazione.

    Ad esempio, invece di usare questa riga di codice: spark.conf.set("fs.azure.account.oauth2.client.secret", "<password>"), usa la seguente riga di codice: spark.conf.set("fs.azure.account.oauth2.client.secret", dbutils.secrets.get(scope = "<scope-name>", key = "<key-name-for-service-credential>")).

    Dopo aver completato questo tutorial, consulta l'articolo Azure Data Lake Storage sul sito Azure Databricks per vedere esempi di questo approccio.

  2. Premi SHIFT + ENTER per eseguire il codice in questo blocco.

  3. Copia e incolla il blocco di codice successivo in una cella diversa, poi premi SHIFT + ENTER per eseguire il codice in questo blocco.

    from pyspark.sql.types import StructType, StructField, DoubleType, IntegerType, StringType
    
    inputSchema = StructType([
    StructField("InvoiceNo", IntegerType(), True),
    StructField("StockCode", StringType(), True),
    StructField("Description", StringType(), True),
    StructField("Quantity", IntegerType(), True),
    StructField("InvoiceDate", StringType(), True),
    StructField("UnitPrice", DoubleType(), True),
    StructField("CustomerID", IntegerType(), True),
    StructField("Country", StringType(), True)
    ])
    
    rawDataDF = (spark.read
     .option("header", "true")
     .schema(inputSchema)
     .csv(adlsPath + 'input')
    )
    
    (rawDataDF.write
      .mode("overwrite")
      .format("delta")
      .saveAsTable("customer_data", path=customerTablePath))
    

    Questo codice crea la tabella Delta di Databricks nel tuo account di archiviazione e poi carica alcuni dati iniziali dal file CSV che hai caricato in precedenza.

  4. Dopo che questo blocco di codice è stato eseguito con successo, rimuovi questo blocco dal tuo notebook.

Aggiungere codice per l'inserimento di righe nella tabella di Databricks Delta

  1. Copiare e incollare il blocco di codice seguente in una cella diversa, ma non eseguire ancora la cella.

    upsertDataDF = (spark
      .read
      .option("header", "true")
      .csv(inputPath)
    )
    upsertDataDF.createOrReplaceTempView("customer_data_to_upsert")
    

    Questo codice inserisce i dati in una vista di tabella temporanea utilizzando dati da un file CSV. Il percorso verso quel file CSV proviene dal widget di input che hai creato in un passaggio precedente.

  2. Copiare e incollare il blocco di codice seguente in una cella diversa. Questo codice unisce il contenuto della vista tabella temporanea alla tabella Databricks Delta.

    %sql
    MERGE INTO customer_data cd
    USING customer_data_to_upsert cu
    ON cd.CustomerID = cu.CustomerID
    WHEN MATCHED THEN
      UPDATE SET
        cd.StockCode = cu.StockCode,
        cd.Description = cu.Description,
        cd.InvoiceNo = cu.InvoiceNo,
        cd.Quantity = cu.Quantity,
        cd.InvoiceDate = cu.InvoiceDate,
        cd.UnitPrice = cu.UnitPrice,
        cd.Country = cu.Country
    WHEN NOT MATCHED
      THEN INSERT (InvoiceNo, StockCode, Description, Quantity, InvoiceDate, UnitPrice, CustomerID, Country)
      VALUES (
        cu.InvoiceNo,
        cu.StockCode,
        cu.Description,
        cu.Quantity,
        cu.InvoiceDate,
        cu.UnitPrice,
        cu.CustomerID,
        cu.Country)
    

Crea un lavoro

Creare un processo per l'esecuzione del notebook creato in precedenza. Successivamente, crei una funzione Azure che esegue questo lavoro quando viene generato un evento.

  1. Seleziona nuovo>lavoro.

  2. Dai un nome al lavoro, scegli il quaderno che hai creato e seleziona un cluster. Quindi, seleziona Crea per creare il processo.

    Il nuovo lavoro appare nella lista Lavori insieme al taccuino e al cluster che hai selezionato.

Crea una Funzione di Azure

Crea una funzione Azure che esegue il lavoro.

  1. Nel tuo workspace Azure Databricks, seleziona il tuo nome utente Azure Databricks nella barra in alto. Dalla lista a tendina, seleziona Impostazioni utente.

  2. Nella scheda Token di accesso selezionare Genera nuovo token.

  3. Copia il token che appare, poi seleziona Fatto.

  4. Nell'angolo superiore dell'area di lavoro di Databricks scegliere l'icona Persone, quindi scegliere Impostazioni utente.

    Screenshot del menu delle impostazioni utente per generare un token di accesso Databricks.

  5. Selezionare il pulsante Genera nuovo token e quindi selezionare il pulsante Genera.

    Assicurati di copiare il token in un luogo sicuro. La tua funzione Azure ha bisogno che questo token si autentichi con Databricks così da poter eseguire il lavoro.

  6. Nel menu del portale di Azure o dalla pagina Home selezionare Crea una risorsa.

  7. Nella pagina Nuovo, selezionare Calcolo>App per le funzioni.

  8. Nella scheda Informazioni di base della pagina Crea app per le funzioni scegliere un gruppo di risorse e quindi modificare o verificare le impostazioni seguenti:

    Impostazione Valore
    Nome della Function App contosoorder
    Stack di runtime .NET
    Pubblicazione Codice
    Sistema operativo Windows
    Tipo di piano Consumo (Serverless)
  9. Seleziona Rivedi e crea e quindi seleziona Crea.

    Al termine della distribuzione, selezionare Vai alla risorsa per aprire la pagina di panoramica dell'app per le funzioni.

  10. Nel gruppo Impostazioni selezionare Configurazione.

  11. Nella pagina Impostazioni applicazione scegliere il pulsante Nuova impostazione applicazione per aggiungere ogni impostazione.

    Schermata dell'aggiunta di una nuova impostazione dell'applicazione nella configurazione della Function App.

    Aggiungi le impostazioni seguenti:

    Nome dell'impostazione Valore
    DBX_INSTANCE Area dell'area di lavoro di Databricks. Ad esempio: westus2.azuredatabricks.net
    DBX_PAT Token di accesso personale generato in precedenza.
    DBX_JOB_ID Identificatore del processo in esecuzione.
  12. Selezionare Salva per eseguire il commit di queste impostazioni.

  13. Nel gruppo Funzioni selezionare Funzioni e quindi crea.

  14. Scegli Griglia di eventi di Azure Trigger.

    Installare l'estensione Microsoft.Azure.WebJobs.Extensions.EventGrid, se viene chiesto di farlo. Se devi installarlo, seleziona di nuovo Griglia di eventi di Azure Trigger per creare la funzione.

    Viene visualizzato il riquadro Nuova funzione.

  15. In Nuova Funzione, inserisci UpsertOrder il nome della funzione e poi seleziona Crea.

  16. Sostituisci il contenuto del file di codice con il seguente codice, poi seleziona Salva:

      #r "Azure.Messaging.EventGrid"
      #r "System.Memory.Data"
      #r "Newtonsoft.Json"
      #r "System.Text.Json"
      using Azure.Messaging.EventGrid;
      using Azure.Messaging.EventGrid.SystemEvents;
      using Newtonsoft.Json;
      using Newtonsoft.Json.Linq;
    
      private static HttpClient httpClient = new HttpClient();
    
      public static async Task Run(EventGridEvent eventGridEvent, ILogger log)
      {
         log.LogInformation("Event Subject: " + eventGridEvent.Subject);
         log.LogInformation("Event Topic: " + eventGridEvent.Topic);
         log.LogInformation("Event Type: " + eventGridEvent.EventType);
         log.LogInformation(eventGridEvent.Data.ToString());
    
         if (eventGridEvent.EventType == "Microsoft.Storage.BlobCreated" || eventGridEvent.EventType == "Microsoft.Storage.FileRenamed") {
            StorageBlobCreatedEventData fileData = eventGridEvent.Data.ToObjectFromJson<StorageBlobCreatedEventData>();
            if (fileData.Api == "FlushWithClose") {
                  log.LogInformation("Triggering Databricks Job for file: " + fileData.Url);
                  var fileUrl = new Uri(fileData.Url);
                  var httpRequestMessage = new HttpRequestMessage {
                     Method = HttpMethod.Post,
                     RequestUri = new Uri(String.Format("https://{0}/api/2.0/jobs/run-now", System.Environment.GetEnvironmentVariable("DBX_INSTANCE", EnvironmentVariableTarget.Process))),
                     Headers = { 
                        { System.Net.HttpRequestHeader.Authorization.ToString(), "Bearer " + System.Environment.GetEnvironmentVariable("DBX_PAT", EnvironmentVariableTarget.Process)},
                        { System.Net.HttpRequestHeader.ContentType.ToString(), "application/json" }
                     },
                     Content = new StringContent(JsonConvert.SerializeObject(new {
                        job_id = System.Environment.GetEnvironmentVariable("DBX_JOB_ID", EnvironmentVariableTarget.Process),
                        notebook_params = new {
                              source_file = String.Join("", fileUrl.Segments.Skip(2))
                        }
                     }))
                  };
                  var response = await httpClient.SendAsync(httpRequestMessage);
                  response.EnsureSuccessStatusCode();
            }
         }
      }
    

    Questo codice analizza le informazioni sull'evento di memorizzazione che è stato generato, e poi crea un messaggio di richiesta con l'URL del file che ha attivato l'evento. Come parte del messaggio, la funzione passa un valore al widget source_file creato in precedenza. Il codice della funzione invia il messaggio al job Databricks e utilizza il token che hai ottenuto in precedenza per l'autenticazione.

Creare una sottoscrizione di Griglia di eventi

In questa sezione, crei un abbonamento a Event Grid che chiama la funzione Azure quando i file vengono caricati sull'account di archiviazione.

  1. Selezionare Integrazione. Nella pagina Integrazione , seleziona Trigger della Griglia Eventi.

  2. Nel riquadro Modifica trigger assegnare all'evento eventGridEventil nome e quindi selezionare Crea sottoscrizione di eventi.

    Nota

    Il nome eventGridEvent corrisponde al nome del parametro che riceve la funzione Azure.

  3. Nella scheda Informazioni di base della pagina Crea sottoscrizione di eventi modificare o verificare le impostazioni seguenti:

    Impostazione Valore
    Name contoso-order-event-subscription
    Tipo di argomento Account di archiviazione
    Risorsa di origine contosoorders
    Nome dell'argomento di sistema <create any name>
    Filtra per tipi di evento BLOB creato e BLOB eliminato
  4. Fare clic su Crea.

Testare la sottoscrizione di Event Grid

  1. Creare un file denominato customer-order.csv, incollare le informazioni seguenti nel file e salvarlo nel computer locale.

    InvoiceNo,StockCode,Description,Quantity,InvoiceDate,UnitPrice,CustomerID,Country
    536371,99999,EverGlow Single,228,1/1/2018 9:01,33.85,20993,Sierra Leone
    
  2. Nel browser di archiviazione, carica questo file nella cartella di input del tuo account di archiviazione.

    Quando carichi un file, viene generato l'evento Microsoft.Storage.BlobCreated. Griglia di eventi invia una notifica a tutti i sottoscrittori dell'evento. In questo caso, la funzione Azure è l'unico abbonato. La funzione di Azure analizza i parametri dell'evento per determinare l'evento che si è verificato. Quindi trasmette l'URL del file al processo Databricks. Il processo Databricks legge il file e aggiunge una riga alla tabella Databricks Delta situata nel tuo account di archiviazione.

  3. Per verificare se il processo è riuscito, visualizzarne le esecuzioni. Visualizzi lo stato di completamento. Per maggiori informazioni su come visualizzare le esecuzioni di un processo, consulta Visualizzare le esecuzioni di un processo.

  4. In una nuova cella del notebook, esegui questa query per visualizzare la tabella Delta aggiornata.

    %sql select * from customer_data
    

    La tabella restituita mostra il record più recente.

    Schermata della query sulla tabella di Databricks Delta che mostra il record più recente.

  5. Per aggiornare questo record, creare un file denominato customer-order-update.csv, incollare le informazioni seguenti nel file e salvarlo nel computer locale.

    InvoiceNo,StockCode,Description,Quantity,InvoiceDate,UnitPrice,CustomerID,Country
    536371,99999,EverGlow Single,22,1/1/2018 9:01,33.85,20993,Sierra Leone
    

    Questo file CSV è quasi identico al precedente, tranne che la quantità dell'ordine viene modificata da 228 a 22.

  6. Nel browser di archiviazione, carica questo file nella cartella di input del tuo account di archiviazione.

  7. Eseguire di nuovo la query select per visualizzare la tabella Delta aggiornata.

    %sql select * from customer_data
    

    La tabella restituita mostra il record aggiornato.

    Schermata della query sulla tabella di Databricks Delta che mostra il record aggiornato.

Pulire le risorse

Quando le risorse non sono più necessarie, eliminare il gruppo di risorse e tutte le risorse correlate. Per eliminare il gruppo risorse, seleziona il gruppo risorse per l'account di archiviazione e seleziona Elimina.

Passaggi successivi