Tutorial: Implementieren des Data Lake-Erfassungsmusters zum Aktualisieren einer Databricks Delta-Tabelle

In diesem Tutorial wird veranschaulicht, wie Sie Ereignisse in einem Speicherkonto mit einem hierarchischen Namespace verarbeiten.

Sie erstellen eine kleine Lösung, mit der Sie eine Databricks-Delta-Tabelle durch Hochladen einer CSV-Datei befüllen können, die einen Kundenauftrag beschreibt. Du baust diese Lösung, indem du ein Event Grid-Abonnement, eine Azure-Funktion und einen Job in Azure Databricks verbindest.

In diesem Tutorial lernen Sie:

  • Erstellen Sie ein Event Grid-Abonnement, das eine Azure-Funktion aufruft.
  • Erstellen Sie eine Azure-Funktion, die eine Benachrichtigung von einem Ereignis empfängt und den Auftrag dann in Azure Databricks ausführt.
  • Erstellen Sie einen Databricks-Auftrag, mit dem eine Kundenbestellung in eine Databricks Delta-Tabelle eingefügt wird, die unter dem Speicherkonto vorhanden ist.

Du baust diese Lösung in umgekehrter Reihenfolge, beginnend mit dem Azure Databricks Workspace.

Voraussetzungen

  • Erstellen eines Speicherkontos mit einem hierarchischen Namespace (Azure Data Lake Storage) In diesem Tutorial wird ein Speicherkonto mit dem Namen contosoorders verwendet.

    Weitere Informationen finden Sie unter Erstellen eines Speicherkontos für die Verwendung mit Azure Data Lake Storage.

  • Vergewissern Sie sich, dass Ihrem Benutzerkonto die Rolle Mitwirkender an Storage-Blobdaten zugewiesen ist.

  • Erstellen Sie einen Dienstprinzipal, erstellen Sie einen geheimen Clientschlüssel, und gewähren Sie dem Dienstprinzipal dann Zugriff auf das Speicherkonto.

    Weitere Informationen dazu finden Sie im Tutorial: Herstellen einer Verbindung mit Azure Data Lake Storage (Schritte 1 bis 3). Nachdem Sie diese Schritte abgeschlossen haben, fügen Sie unbedingt die Werte für Tenant-ID, App-ID und Client Secret in eine Textdatei ein. Du brauchst diese Werte bald.

  • Wenn Sie kein Azure-Abonnement besitzen, können Sie ein kostenloses Konto erstellen, bevor Sie beginnen.

Erstellen eines Verkaufsauftrags

Zuerst erstellen Sie eine CSV-Datei, die eine Verkaufsbestellung beschreibt, und laden Sie diese Datei dann auf das Speicherkonto hoch. Später verwenden Sie die Daten aus dieser Datei, um die erste Zeile in Ihrer Databricks-Delta-Tabelle auszufüllen.

  1. Navigieren Sie im Azure-Portal zu Ihrem neuen Speicherkonto.

  2. Wählen Sie Speicherbrowser>Blobcontainer>Container hinzufügen aus und erstellen Sie einen neuen Container namens data.

    Screenshot des Erstellens eines Containers im Azure Storage-Browser.

  3. Erstellen Sie im Container data ein Verzeichnis mit dem Namen input.

  4. Fügen Sie in einem Text-Editor den folgenden Text ein:

    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. Speichere diese Datei auf deinem lokalen Computer und nenne sie data.csv.

  6. Im Speicherbrowser laden Sie diese Datei in den Eingabeordner hoch.

Erstellen eines Auftrags in Azure Databricks

In diesem Abschnitt führen Sie folgende Aufgaben aus:

  • Erstellen eines Azure Databricks-Arbeitsbereichs
  • Erstellen Sie ein Notebook.
  • Erstellen und Auffüllen einer Databricks Delta-Tabelle
  • Hinzufügen von Code, mit dem Zeilen in die Databricks Delta-Tabelle eingefügt werden
  • Erstellen Sie einen Job.

Erstellen eines Azure Databricks-Arbeitsbereichs

In diesem Abschnitt erstellen Sie einen Azure Databricks-Arbeitsbereich mithilfe des Azure-Portals.

  1. Erstellen eines Azure Databricks-Arbeitsbereichs Nenne den Arbeitsbereich contoso-orders. Weitere Informationen finden Sie unter Erstellen eines Azure Databricks-Arbeitsbereichs.

  2. Erstellen eines Clusters Nennen Sie den Cluster customer-order-cluster. Weitere Informationen finden Sie unter Erstellen eines Clusters.

  3. Erstellen Sie ein Notebook. Nennen Sie das Notebook configure-customer-table, und wählen Sie Python als Standardsprache für das Notebook aus. Weitere Informationen finden Sie unter Erstellen eines Notebooks.

Erstellen und Auffüllen einer Databricks Delta-Tabelle

  1. Kopieren Sie im erstellten Notebook den folgenden Codeblock, und fügen Sie ihn in die erste Zelle ein, aber führen Sie den Code noch nicht aus.

    Ersetzen Sie die Platzhalterwerte appId, password und tenant in diesem Codeblock durch die Werte, die Sie beim Abschließen der Voraussetzungen dieses Tutorials erfasst haben.

    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'
    

    Mit diesem Code wird ein Widget mit dem Namen source_file erstellt. Später erstellen Sie eine Azure-Funktion, mit der dieser Code aufgerufen und ein Dateipfad an das Widget übergeben wird. Dieser Code authentifiziert außerdem Ihren Service-Principal beim Speicherkonto und erstellt einige Variablen, die Sie in anderen Zellen verwenden.

    Hinweis

    In einer Produktionsumgebung empfiehlt es sich, Ihren Authentifizierungsschlüssel in Azure Databricks zu speichern. Füge dann deinem Codeblock statt des Authentifizierungsschlüssels einen Lookup-Schlüssel hinzu.

    Zum Beispiel verwenden Sie statt dieser Codezeile: spark.conf.set("fs.azure.account.oauth2.client.secret", "<password>"), die folgende Codezeile: spark.conf.set("fs.azure.account.oauth2.client.secret", dbutils.secrets.get(scope = "<scope-name>", key = "<key-name-for-service-credential>")).

    Nachdem Sie dieses Tutorial abgeschlossen haben, sehen Sie sich den Artikel über Azure Data Lake Storage auf der Azure Databricks-Website an, um Beispiele für diesen Ansatz zu sehen.

  2. Drücken Sie SHIFT + ENTER, um den Code in diesem Block auszuführen.

  3. Kopieren und fügen Sie den folgenden Codeblock in eine andere Zelle ein und drücken Sie dann SHIFT + ENTER , um den Code in diesem Block auszuführen.

    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))
    

    Dieser Code erstellt die Databricks-Delta-Tabelle in deinem Speicherkonto und lädt dann einige Anfangsdaten aus der zuvor hochgeladenen CSV-Datei.

  4. Nachdem dieser Codeblock erfolgreich ausgeführt wurde, entferne diesen Codeblock aus deinem Notizbuch.

Hinzufügen von Code, mit dem Zeilen in die Databricks Delta-Tabelle eingefügt werden

  1. Kopieren Sie den folgenden Codeblock, und fügen Sie ihn in eine andere Zelle ein, aber führen Sie diese Zelle nicht aus.

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

    Dieser Code fügt Daten in eine temporäre Tabellenansicht ein, indem Daten aus einer CSV-Datei verwendet werden. Der Pfad zu dieser CSV-Datei stammt vom Eingabe-Widget, das du in einem früheren Schritt erstellt hast.

  2. Kopieren Sie den folgenden Codeblock, und fügen Sie ihn in eine andere Zelle ein. Dieser Code führt den Inhalt der temporären Tabellenansicht mit der Databricks Delta-Tabelle zusammen.

    %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)
    

Einen Job erstellen

Erstelle einen Job, der das Notebook ausführt, das du zuvor erstellt hast. Später erstellst du eine Azure-Funktion, die diesen Job ausführt, wenn ein Ereignis ausgelöst wird.

  1. NeueStelle auswählen>.

  2. Gib dem Job einen Namen, wähle das Notebook, das du erstellt hast, und wähle einen Cluster aus. Wählen Sie anschließend Erstellen aus, um den Auftrag zu erstellen.

    Der neue Job erscheint in der Jobs-Liste zusammen mit dem von dir ausgewählten Notizbuch und Cluster.

Erstellen einer Azure Function

Erstellen Sie eine Azure-Funktion, die den Job ausführt.

  1. Wählen Sie in Ihrem Azure Databricks Workspace Ihren Azure Databricks-Benutzernamen in der oberen Leiste aus. Wählen Sie in der Dropdown-Liste Benutzereinstellungen aus.

  2. Wählen Sie auf der Registerkarte Zugriffstoken die Option Neues Token generieren aus.

  3. Kopiere das angezeigte Token und wähle dann Fertig aus.

  4. Wählen Sie in der oberen Ecke des Databricks-Arbeitsbereichs das Symbol „Personen“, und wählen Sie dann Benutzereinstellungen.

    Screenshot des Benutzereinstellungsmenüs zur Generierung eines Databricks-Zugriffstokens.

  5. Wählen Sie die Schaltfläche Neues Token generieren und dann die Schaltfläche Generieren aus.

    Stelle sicher, dass du den Token an einen sicheren Ort kopierst. Deine Azure-Funktion benötigt dieses Token, um sich mit Databricks zu authentifizieren und den Job ausführen zu können.

  6. Wählen Sie im Menü des Azure-Portals oder auf der Startseite die Option Ressource erstellen aus.

  7. Wählen Sie auf der Seite Neu die Option Compute>Function App aus.

  8. Wählen Sie auf der Registerkarte Grundlagen der Seite Funktions-App erstellen eine Ressourcengruppe aus, und ändern bzw. überprüfen Sie die folgenden Einstellungen:

    Einstellung Wert
    Name der Funktions-App contosoorder
    Runtimestapel .NET
    Veröffentlichen Code
    Betriebssystem Windows
    Plantyp Verbrauch (Serverless)
  9. Klicken Sie auf Überprüfen und erstellen und dann auf Erstellen.

    Wählen Sie nach Abschluss der Bereitstellung die Option Zu Ressource wechseln aus, um die Übersichtsseite der Funktions-App zu öffnen.

  10. Wählen Sie in der Gruppe Einstellungen die Option Konfiguration aus.

  11. Wählen Sie auf der Seite Anwendungseinstellungen die Schaltfläche Neue Anwendungseinstellung, um die einzelnen Einstellungen hinzuzufügen.

    Screenshot zum Hinzufügen einer neuen Anwendungseinstellung in der Konfiguration der Function App.

    Fügen Sie die folgenden Einstellungen hinzu:

    Einstellungsname Wert
    DBX_INSTANCE Die Region Ihres Databricks-Arbeitsbereichs. Beispiel: westus2.azuredatabricks.net
    DBX_PAT Das persönliche Zugriffstoken, das Sie zuvor generiert haben.
    DBX_JOB_ID Der Bezeichner des ausgeführten Auftrags.
  12. Wählen Sie Speichern aus, um diese Einstellungen zu committen.

  13. Wählen Sie in der Gruppe Funktionen die Option Funktionen und anschließend Erstellen aus.

  14. Wählen Sie Azure Event Grid-Trigger.

    Installieren Sie die Erweiterung Microsoft.Azure.WebJobs.Extensions.EventGrid, wenn Sie dazu aufgefordert werden. Wenn du es installieren musst, wähle erneut Azure Event Grid Trigger, um die Funktion zu erstellen.

    Der Bereich Neue Funktion wird angezeigt.

  15. In Neue Funktion geben Sie den Funktionsnamen ein UpsertOrder und wählen Sie dann Erstellen.

  16. Ersetzen Sie den Inhalt der Codedatei durch folgenden Code und wählen Sie dann Speichern:

      #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();
            }
         }
      }
    

    Dieser Code parselt Informationen über das ausgelöste Speicherereignis und erstellt dann eine Anfrage mit der URL der Datei, die das Ereignis ausgelöst hat. Im Rahmen der Nachricht übergibt die Funktion einen Wert an das Widget source_file, das Sie zuvor erstellt haben. Der Funktionscode sendet die Nachricht an den Databricks-Job und verwendet das zuvor erhaltene Token als Authentifizierung.

Erstellen eines Event Grid-Abonnements

In diesem Abschnitt erstellen Sie ein Event Grid-Abonnement, das die Azure-Funktion aufruft, wenn Dateien auf das Speicherkonto hochgeladen werden.

  1. Wählen Sie Integration aus. Auf der Integrationsseite wählen Sie Event Grid Trigger.

  2. Legen Sie im Bereich Trigger bearbeiten den Namen des Ereignisses auf eventGridEvent fest, und wählen Sie anschließend Ereignisabonnement erstellen aus.

    Hinweis

    Der Name eventGridEvent entspricht dem Parameternamen, den die Azure-Funktion empfängt.

  3. Ändern bzw. überprüfen Sie auf der Registerkarte Grundlagen der Seite Ereignisabonnement erstellen die folgenden Einstellungen:

    Einstellung Wert
    Name contoso-order-event-subscription
    Thementyp Speicherkonto
    Quellressource contosoorders
    Name des Systemthemas <create any name>
    Nach Ereignistypen filtern „Blob erstellt“ und „Blob gelöscht“
  4. Wählen Sie "Erstellen" aus.

Testen des Event Grid-Abonnements

  1. Erstellen Sie eine Datei mit dem Namen customer-order.csv, fügen Sie die folgenden Informationen in diese Datei ein, und speichern Sie sie auf Ihrem lokalen Computer.

    InvoiceNo,StockCode,Description,Quantity,InvoiceDate,UnitPrice,CustomerID,Country
    536371,99999,EverGlow Single,228,1/1/2018 9:01,33.85,20993,Sierra Leone
    
  2. Im Speicherbrowser laden Sie diese Datei in den Eingabeordner Ihres Speicherkontos hoch.

    Wenn du eine Datei hochlädst, wird das Microsoft.Storage.BlobCreated-Ereignis ausgelöst. Event Grid benachrichtigt alle Abonnenten über dieses Ereignis. In diesem Fall ist die Azure-Funktion der einzige Abonnent. Die Azure-Funktion analysiert die Ereignisparameter, um zu ermitteln, welches Ereignis eingetreten ist. Anschließend wird die URL der Datei an den Databricks-Job weitergegeben. Der Databricks-Job liest die Datei und fügt der Databricks-Delta-Tabelle eine Zeile hinzu, die sich in Ihrem Speicherkonto befindet.

  3. Um zu überprüfen, ob der Auftrag erfolgreich war, sehen Sie sich die Ausführungen für Ihren Auftrag an. Du siehst den Abschlussstatus. Weitere Informationen dazu, wie Sie Ausführungen für einen Job anzeigen, finden Sie unter Ausführungen für einen Job anzeigen.

  4. Führe in einer neuen Arbeitsbuchzelle diese Abfrage aus, um die aktualisierte Delta-Tabelle zu sehen.

    %sql select * from customer_data
    

    Die zurückgegebene Tabelle zeigt den neuesten Datensatz.

    Screenshot der Abfrage einer Databricks-Delta-Tabelle mit dem neuesten Datensatz.

  5. Erstellen Sie zum Aktualisieren dieses Datensatzes eine Datei mit dem Namen customer-order-update.csv, fügen Sie die folgenden Informationen in diese Datei ein, und speichern Sie sie auf Ihrem lokalen Computer.

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

    Diese CSV-Datei ist fast identisch mit der vorherigen, außer dass die Bestellmenge von 228 zu geändert wird 22.

  6. Im Speicherbrowser laden Sie diese Datei in den Eingabeordner Ihres Speicherkontos hoch.

  7. Führen Sie die Abfrage select erneut aus, um die aktualisierte Delta-Tabelle anzuzeigen.

    %sql select * from customer_data
    

    In der zurückgegebenen Tabelle wird der aktualisierte Datensatz angezeigt.

    Screenshot der Databricks-Delta-Tabellenanfrage zeigt den aktualisierten Datensatz.

Bereinigen von Ressourcen

Wenn Sie die Ressourcen nicht mehr benötigen, löschen Sie die Ressourcengruppe und alle zugehörigen Ressourcen. Um die Ressourcengruppe zu löschen, wählen Sie die Ressourcengruppe für das Speicherkonto und wählen Sie Löschen.

Nächste Schritte

Reacting to Blob storage events (preview) (Reagieren auf Blob Storage-Ereignisse (Vorschauversion))