Tutoriel : Implémenter le modèle de capture de lac de données pour mettre à jour une table Delta Databricks

Ce tutoriel vous montre comment gérer des événements dans un compte de stockage doté d’un espace de noms hiérarchique.

Vous construisez une petite solution qui vous permet de remplir une table Delta Databricks en téléchargeant un fichier de valeurs séparées par virgules (CSV) qui décrit une commande de vente. Vous construisez cette solution en connectant un abonnement Event Grid, une fonction Azure et un Job dans Azure Databricks.

Dans ce tutoriel, vous allez :

  • Créer un abonnement Event Grid qui appelle une fonction Azure
  • Créer une fonction Azure qui reçoit une notification d’un événement, puis exécute le travail dans Azure Databricks
  • Créer une tâche Databricks qui insère une commande client dans une table Databricks Delta située dans le compte de stockage.

Vous construisez cette solution dans l’ordre inverse, en commençant par l’espace de travail Azure Databricks.

Prérequis

Créer une commande client

D’abord, créez un fichier CSV qui décrit une commande de vente, puis téléchargez ce fichier sur le compte de stockage. Plus tard, vous utilisez les données de ce fichier pour remplir la première ligne de votre table Delta Databricks.

  1. Accédez à votre nouveau compte de stockage dans le portail Azure.

  2. Sélectionnez Navigateur de stockage>Conteneurs blob>Ajouter un conteneur et créez un nouveau conteneur nommé data.

    Capture d’écran de la création d’un conteneur dans le navigateur stockage Azure.

  3. Dans le conteneur data, créez un répertoire nommé input.

  4. Dans un éditeur de texte, collez le texte suivant.

    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. Enregistrez ce fichier sur votre ordinateur local et donnez-lui un nomdata.csv.

  6. Dans le navigateur de stockage, téléchargez ce fichier dans le dossier d’entrée .

Créer un travail dans Azure Databricks

Dans cette section, vous effectuez ces tâches :

  • Créer un espace de travail Azure Databricks.
  • Créez un bloc-notes.
  • Créer et remplir une table Databricks Delta
  • Ajouter du code qui insère des lignes dans la table Databricks Delta
  • Créez un travail.

Créer un espace de travail Azure Databricks

Dans cette section, vous créez un espace de travail Azure Databricks en utilisant le portail Azure.

  1. Créer un espace de travail Azure Databricks. Nomme l’espace de travail contoso-orders. Consultez Créer un espace de travail Azure Databricks.

  2. Créer un cluster. Nommez le cluster customer-order-cluster. Voir Créez un cluster.

  3. Créez un bloc-notes. Nommez le notebook configure-customer-table et choisissez Python comme langage par défaut du notebook. Consultez Création d’un notebook.

Créer et remplir une table Databricks Delta

  1. Dans le notebook que vous avez créé, copiez et collez le bloc de code suivant dans la première cellule, mais n’exécutez pas ce code pour l’instant.

    Remplacez les valeurs d’espace réservé appId, password et tenant dans ce bloc de code par les valeurs que vous avez collectées lors de la réalisation des prérequis de ce tutoriel.

    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'
    

    Ce code crée un widget nommé source_file. Vous créerez plus tard une fonction Azure qui appelle ce code et transmet un chemin de fichier à ce widget. Ce code authentifie également votre principal de service avec le compte de stockage, et crée certaines variables que vous utilisez dans d’autres cellules.

    Note

    Dans un environnement de production, songez à stocker votre clé d’authentification dans Azure Databricks. Ensuite, ajoutez une clé de recherche à votre bloc de code au lieu de la clé d’authentification.

    Par exemple, au lieu d’utiliser cette ligne de code : spark.conf.set("fs.azure.account.oauth2.client.secret", "<password>"), utilisez la ligne suivante : spark.conf.set("fs.azure.account.oauth2.client.secret", dbutils.secrets.get(scope = "<scope-name>", key = "<key-name-for-service-credential>")).

    Après avoir terminé ce tutoriel, consultez l’article Azure Data Lake Storage sur le site Azure Databricks pour voir des exemples de cette approche.

  2. Appuyez sur SHIFT + ENTER pour exécuter le code dans ce bloc.

  3. Copiez-collez le bloc de code suivant dans une autre cellule, puis appuyez sur SHIFT + ENTER pour exécuter le code de ce bloc.

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

    Ce code crée la table Delta Databricks dans votre compte de stockage, puis charge certaines données initiales du fichier CSV que vous avez téléchargé précédemment.

  4. Après que ce bloc de code ait fonctionné avec succès, retirez-le de votre carnet.

Ajouter du code qui insère des lignes dans la table Databricks Delta

  1. Copiez le bloc de code suivant et collez-le dans une autre cellule, mais n’exécutez pas cette cellule pour l’instant.

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

    Ce code insère des données dans une vue de table temporaire en utilisant des données provenant d’un fichier CSV. Le chemin vers ce fichier CSV provient du widget d’entrée que vous avez créé à une étape précédente.

  2. Copiez et collez le bloc de code suivant dans une autre cellule. Ce code fusionne le contenu de la vue de la table temporaire avec la table 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)
    

Créer un travail

Créez un travail qui exécute le carnet que vous avez créé précédemment. Plus tard, vous créez une fonction Azure qui exécute cette tâche lorsqu’un événement est lancé.

  1. Sélectionnez un nouveau>poste.

  2. Donnez un nom au poste, choisissez le carnet que vous avez créé, puis sélectionnez un cluster. Sélectionnez ensuite Créer pour créer le travail.

    Le nouveau poste apparaît dans la liste des emplois avec le carnet et le cluster que vous avez sélectionnés.

Création d’une fonction Azure

Créez une fonction Azure qui exécute le travail.

  1. Dans votre espace de travail Azure Databricks, sélectionnez votre nom d’utilisateur Azure Databricks dans la barre supérieure. Dans la liste déroulante, sélectionnez Paramètres utilisateur.

  2. Sous l’onglet Jetons d’accès, sélectionnez Générer un nouveau jeton.

  3. Copiez le jeton qui apparaît, puis sélectionnez Terminé.

  4. Dans le coin supérieur de l’espace de travail Databricks, cliquez sur l’icône représentant un personnage, puis sur Paramètres utilisateur.

    Capture d’écran du menu des paramètres utilisateur pour générer un jeton d’accès Databricks.

  5. Sélectionnez le bouton Générer un nouveau jeton, puis le bouton Générer.

    Assurez-vous de copier le jeton dans un endroit sûr. Votre fonction Azure a besoin que ce jeton s’authentifie avec Databricks afin de pouvoir exécuter le travail.

  6. Dans le menu du portail Azure ou dans la page Accueil, sélectionnez Créer une ressource.

  7. Dans la page Nouveau, sélectionnez Calcul>Application de fonction.

  8. Sous l’onglet Informations de base de la page Créer une application de fonction, choisissez un groupe de ressources, puis modifiez ou vérifiez les paramètres suivants :

    Paramètres Valeur
    Nom de l’application de fonction contosoorder
    Pile d’exécution .NET
    Publier Code
    Système d'exploitation Windows
    Type de plan Consommation (Sans serveur)
  9. Sélectionnez Examiner + créer, puis sélectionnez Créer.

    Lorsque le déploiement est terminé, sélectionnez Accéder à la ressource pour ouvrir la page de vue d’ensemble de l’application de fonction.

  10. Dans le groupe Paramètres, sélectionnez Configuration.

  11. Dans la page Paramètres de l’application, cliquez sur le bouton Nouveau paramètre d’application pour ajouter chaque paramètre.

    Capture d’écran de l’ajout d’un nouveau paramètre d’application dans la configuration de l’application Fonction.

    Ajoutez les paramètres suivants :

    Nom du paramètre Valeur
    DBX_INSTANCE La région de votre espace de travail Databricks. Par exemple : westus2.azuredatabricks.net
    DBX_PAT Le jeton d’accès personnel que vous avez généré précédemment.
    DBX_JOB_ID L’identificateur de la tâche en cours d’exécution.
  12. Sélectionnez Enregistrer pour valider ces paramètres.

  13. Dans le groupe Fonctions, sélectionnez Fonctions, puis Créer.

  14. Choisissez le déclencheur Azure Event Grid.

    Si vous y êtes invité, installez l’extension Microsoft.Azure.WebJobs.Extensions.EventGrid. Si vous devez l’installer, sélectionnez à nouveau Azure Event Grid Trigger pour créer la fonction.

    Le volet Nouvelle fonction s’affiche.

  15. Dans Nouvelle fonction, saisissez UpsertOrder le nom de la fonction, puis sélectionnez Créer.

  16. Remplacez le contenu du fichier de code par le code suivant, puis sélectionnez Enregistrer :

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

    Ce code analyse les informations sur l’événement de stockage qui a été lancé, puis crée un message de requête avec l’URL du fichier qui a déclenché l’événement. La fonction transfère une valeur au widget source_file que vous avez créé précédemment dans le message. Le code de fonction envoie le message au travail Databricks et utilise le jeton que vous avez obtenu précédemment comme authentification.

Créer un abonnement Event Grid

Dans cette section, vous créez un abonnement Event Grid qui appelle la fonction Azure lorsque les fichiers sont téléchargés sur le compte de stockage.

  1. Sélectionnez Intégration. Dans la page Intégration , sélectionnez Déclencheur de grille d’événements.

  2. Dans le volet Modifier le déclencheur, nommez l’événement eventGridEvent, puis sélectionnez Créer un abonnement à un événement.

    Note

    Le nom eventGridEvent correspond au nom du paramètre que reçoit la fonction Azure.

  3. Sous l’onglet Informations de base de la page Créer un abonnement à un événement, modifiez ou vérifiez les paramètres suivants :

    Paramètres Valeur
    Nom contoso-order-event-subscription
    Type de rubrique Compte de stockage
    Ressource d’origine contosoorders
    Nom de la rubrique système <create any name>
    Filtrer les types d’événements Blob créé et blob supprimé
  4. Cliquez sur Créer.

Tester l’abonnement Event Grid

  1. Créez un fichier nommé customer-order.csv, collez les informations suivantes dans ce fichier, puis enregistrez-le sur votre ordinateur local.

    InvoiceNo,StockCode,Description,Quantity,InvoiceDate,UnitPrice,CustomerID,Country
    536371,99999,EverGlow Single,228,1/1/2018 9:01,33.85,20993,Sierra Leone
    
  2. Dans le navigateur de stockage, téléchargez ce fichier dans le dossier d’entrée de votre compte de stockage.

    Lorsque vous chargez un fichier, cela déclenche l’événement Microsoft.Storage.BlobCreated. Event Grid informe tous les abonnés de cet événement. Dans ce cas, la fonction Azure est le seul abonné. La fonction Azure analyse les paramètres d’événement pour déterminer l’événement qui s’est produit. Il transmet ensuite l’URL du fichier au job Databricks. La tâche Databricks lit le fichier et ajoute une ligne à la table Databricks Delta située dans votre compte de stockage.

  3. Pour vérifier si la tâche a réussi, consultez les exécutions de votre tâche. Vous voyez un état d’avancement. Pour plus d’informations sur la façon d’afficher les exécutions d’une tâche, voir Afficher les exécutions d’une tâche.

  4. Dans une nouvelle cellule de classeur, lancez cette requête pour voir la table Delta mise à jour.

    %sql select * from customer_data
    

    La table renvoyée affiche l’enregistrement le plus récent.

    Capture d’écran de la requête de la table Delta de Databricks montrant le dernier enregistrement.

  5. Pour mettre à jour cet enregistrement, créez un fichier nommé customer-order-update.csv, collez les informations suivantes dans ce fichier, puis enregistrez-le sur votre ordinateur local.

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

    Ce fichier CSV est presque identique au précédent, sauf que la quantité de l’ordre est modifiée de 228 à 22.

  6. Dans le navigateur de stockage, téléchargez ce fichier dans le dossier d’entrée de votre compte de stockage.

  7. Réexécutez la requête select pour voir la table delta mise à jour.

    %sql select * from customer_data
    

    La table renvoyée affiche l’enregistrement mis à jour.

    Capture d’écran de la requête de la table Delta de Databricks montrant l’enregistrement mis à jour.

Nettoyer les ressources

Lorsque vous n’avez plus besoin des ressources, supprimez le groupe de ressources et toutes les ressources associées. Pour supprimer le groupe de ressources, sélectionnez le groupe de ressources pour le compte de stockage et sélectionnez Supprimer.

Étapes suivantes