Samouczek: implementowanie wzorca przechwytywania usługi Data Lake w celu zaktualizowania tabeli delty usługi Databricks

W tym samouczku pokazano, jak obsługiwać zdarzenia dotyczące konta magazynu z hierarchiczną przestrzenią nazw.

Budujesz małe rozwiązanie, które pozwala wypełnić tabelę Databricks Delta, przesyłając plik wartości rozdzielonych przecinkami (CSV) opisujący zamówienie sprzedaży. To rozwiązanie budujesz, łącząc subskrypcję Event Grid, funkcję Azure oraz pracę w Azure Databricks.

W tym samouczku nauczysz się następujących rzeczy:

  • Utwórz subskrypcję usługi Event Grid, która wywołuje funkcję platformy Azure.
  • Utwórz funkcję platformy Azure, która odbiera powiadomienie z zdarzenia, a następnie uruchamia zadanie w usłudze Azure Databricks.
  • Utwórz zadanie Databricks, które wstawia zamówienie klienta do tabeli Databricks Delta znajdującej się na koncie magazynu danych.

To rozwiązanie budujesz w odwrotnej kolejności, zaczynając od przestrzeni roboczej Azure Databricks.

Wymagania wstępne

Utwórz zamówienie sprzedaży

Najpierw stwórz plik CSV opisujący zamówienie sprzedaży, a następnie prześlij ten plik na konto magazynowe. Później używasz danych z tego pliku, aby wypełnić pierwszy wiersz w tabeli Databricks Delta.

  1. Przejdź do swojego nowego konta magazynowego w portalu Azure.

  2. Wybierz przeglądarkę> magazynowąKontenery> Blob Dodajkontener i stwórz nowy kontener o nazwie data.

    Zrzut ekranu tworzenia kontenera w przeglądarce Azure Storage.

  3. W kontenerze danych utwórz katalog o nazwie input.

  4. Wklej następujący tekst do edytora tekstów.

    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. Zapisz ten plik na swoim komputerze lokalnym i nazwij go data.csv.

  6. W Przeglądarce Storage prześlij ten plik do folderu wejściowego .

Tworzenie zadania w usłudze Azure Databricks

W tej sekcji wykonujesz następujące zadania:

  • Tworzenie obszaru roboczego usługi Azure Databricks.
  • Utwórz notatnik.
  • Utwórz i wypełnij tabelę delty usługi Databricks.
  • Dodaj kod, który wstawia wiersze do tabeli Delta w usłudze Databricks.
  • Utwórz zadanie.

Tworzenie obszaru roboczego usługi Azure Databricks

W tej sekcji tworzysz przestrzeń roboczą Azure Databricks, korzystając z portalu Azure.

  1. Tworzenie obszaru roboczego usługi Azure Databricks. Nazwij przestrzeń contoso-ordersroboczą. Zobacz Tworzenie obszaru roboczego usługi Azure Databricks.

  2. Tworzenie klastra. Nadaj klastrowi customer-order-clusternazwę . Zobacz Tworzenie klastra.

  3. Utwórz notatnik. Nadaj notesowi configure-customer-table nazwę i wybierz język Python jako domyślny język notesu. Zobacz Tworzenie notesu.

Tworzenie i wypełnianie tabeli delty usługi Databricks

  1. W utworzonym notesie skopiuj i wklej następujący blok kodu do pierwszej komórki, ale nie uruchamiaj jeszcze tego kodu.

    Zastąp wartości symboli zastępczych appId, password i tenant w tym bloku kodu wartościami, które zebrano podczas wykonywania czynności wstępnych tego samouczka.

    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'
    

    Ten kod tworzy widżet o nazwie source_file. Później utworzysz funkcję platformy Azure, która wywołuje ten kod i przekazuje ścieżkę pliku do tego widżetu. Kod ten uwierzytelnia również Twój podmiot usługowy za pomocą konta pamięci oraz tworzy zmienne, których używasz w innych komórkach.

    Uwaga

    W środowisku produkcyjnym rozważ przechowywanie klucza uwierzytelniania w usłudze Azure Databricks. Następnie dodaj klucz wyszukiwania do bloku kodu zamiast klucza uwierzytelniającego.

    Na przykład, zamiast używać tej linii kodu: spark.conf.set("fs.azure.account.oauth2.client.secret", "<password>"), użyj następującej linii kodu: spark.conf.set("fs.azure.account.oauth2.client.secret", dbutils.secrets.get(scope = "<scope-name>", key = "<key-name-for-service-credential>")).

    Po ukończeniu tego samouczka zobacz artykuł Azure Data Lake Storage na stronie Azure Databricks, aby zobaczyć przykłady tego podejścia.

  2. Naciśnij SHIFT + ENTER , aby uruchomić kod w tym bloku.

  3. Skopiuj i wklej następujący blok kodu do innej komórki, a następnie naciśnij SHIFT + ENTER , aby uruchomić kod w tym bloku.

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

    Ten kod tworzy tabelę Databricks Delta w Twoim koncie pamięci masowej, a następnie ładuje początkowe dane z pliku CSV, który wcześniej przesłałeś.

  4. Po pomyślnym uruchomieniu tego bloku kodu usuń go ze swojego notesu.

Dodaj kod umożliwiający wstawianie wierszy do tabeli Delta w usłudze Databricks

  1. Skopiuj i wklej następujący blok kodu do innej komórki, ale nie uruchamiaj tej komórki.

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

    Kod ten wstawia dane do tymczasowego widoku tabelowego, wykorzystując dane z pliku CSV. Ścieżka do tego pliku CSV pochodzi z widżetu wejściowego, który utworzyłeś w wcześniejszym kroku.

  2. Skopiuj i wklej następujący blok kodu do innej komórki. Ten kod scala zawartość tymczasowego widoku tabeli z tabelą Delta w usłudze Databricks.

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

Utwórz pracę

Stwórz zadanie, które uruchamia notatnik, który wcześniej stworzyłeś. Później tworzysz funkcję Azure, która uruchamia to zadanie po wywołaniu zdarzenia.

  1. Wybierz Nową>pracę.

  2. Nadaj zadaniu nazwę, wybierz notatnik, który stworzyłeś, i wybierz klaster. Następnie wybierz pozycję Utwórz , aby utworzyć zadanie.

    Nowe zadanie pojawia się na liście Jobs wraz z wybranym przez Ciebie notatnikiem i klastrem.

Tworzenie funkcji platformy Azure

Stwórz funkcję Azure, która uruchamia to zadanie.

  1. W swojej przestrzeni roboczej Azure Databricks wybierz nazwę użytkownika Azure Databricks w górnym pasku. Z listy rozwijanej wybierz Ustawienia użytkownika.

  2. Na karcie Tokeny dostępu wybierz pozycję Generuj nowy token.

  3. Skopiuj token, który się pojawi, a następnie wybierz Gotowe.

  4. W górnym rogu obszaru roboczego usługi Databricks wybierz ikonę osoby, a następnie wybierz pozycję Ustawienia użytkownika.

    Zrzut ekranu menu ustawień użytkownika do generowania tokena dostępu Databricks.

  5. Wybierz przycisk Generuj nowy token, a następnie wybierz przycisk Generuj.

    Upewnij się, że skopiujesz token w bezpieczne miejsce. Twoja funkcja Azure potrzebuje tego tokena do uwierzytelniania się z Databricks, aby móc wykonać zadanie.

  6. W menu witryny Azure Portal lub na stronie głównej wybierz pozycję Utwórz zasób.

  7. Na stronie Nowy wybierz pozycję Obliczenia>Aplikacja funkcji.

  8. Na stronie Tworzenie aplikacji funkcji na karcie Podstawy wybierz grupę zasobów, a następnie zmień lub sprawdź następujące ustawienia:

    Ustawienie Wartość
    Nazwa aplikacji funkcji contosoorder
    Stos środowiska uruchomieniowego .NET
    Publikowanie Code
    System operacyjny Windows
    Typ planu Zużycie (bezserwerowe)
  9. Wybierz opcję Przejrzyj i utwórz, a następnie wybierz pozycję Utwórz.

    Po zakończeniu wdrażania wybierz pozycję Przejdź do zasobu , aby otworzyć stronę przeglądu aplikacji funkcji.

  10. W grupie Ustawienia wybierz pozycję Konfiguracja.

  11. Na stronie Ustawienia aplikacji wybierz przycisk Nowe ustawienie aplikacji, aby dodać każde ustawienie.

    Zrzut ekranu dodawania nowego ustawienia aplikacji w konfiguracji aplikacji Function.

    Dodaj następujące ustawienia:

    Nazwa ustawienia Wartość
    DBX_INSTANCE Region obszaru roboczego usługi Databricks. Na przykład: westus2.azuredatabricks.net.
    DBX_PAT Osobisty token dostępu wygenerowany wcześniej.
    DBX_JOB_ID Identyfikator uruchomionego zadania.
  12. Wybierz pozycję Zapisz , aby zatwierdzić te ustawienia.

  13. W grupie Funkcje wybierz pozycję Funkcje, a następnie wybierz pozycję Utwórz.

  14. Wybierz pozycję Wyzwalacz usługi Azure Event Grid.

    Zainstaluj rozszerzenie Microsoft.Azure.WebJobs.Extensions.EventGrid, jeśli zostanie wyświetlony monit o to. Jeśli musisz ją zainstalować, wybierz ponownie Azure Event Grid Trigger, aby utworzyć tę funkcję.

    Zostanie wyświetlony panel Nowa funkcja.

  15. W Nowa funkcja wpisz UpsertOrder nazwę funkcji, a następnie wybierz Create.

  16. Zastąp zawartość pliku kodu następującym kodem, a następnie wybierz Zapisz:

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

    Kod ten analizuje informacje o zdarzeniu pamięci, które zostało wywołane, a następnie tworzy wiadomość żądania z adresem URL pliku, który wywołał to zdarzenie. W ramach komunikatu funkcja przekazuje wartość do utworzonego wcześniej widżetu source_file . Kod funkcji wysyła komunikat do zadania Databricks i używa tokena uzyskanego wcześniej do uwierzytelniania.

Tworzenie subskrypcji usługi Event Grid

W tej sekcji tworzysz subskrypcję Event Grid, która wywołuje funkcję Azure podczas przesyłania plików na konto pamięci.

  1. Wybierz pozycję Integracja. Na stronie Integracja wybierz Wyzwalacz Siatki Zdarzeń.

  2. W okienku Edytowanie wyzwalacza nadaj zdarzeniu eventGridEventnazwę , a następnie wybierz pozycję Utwórz subskrypcję zdarzeń.

    Uwaga

    Nazwa eventGridEvent odpowiada nazwie parametru, którą otrzymuje funkcja Azure.

  3. Na karcie Podstawowe na stronie Tworzenie subskrypcji zdarzeń zmień lub sprawdź następujące ustawienia:

    Ustawienie Wartość
    Nazwisko contoso-order-event-subscription
    Typ tematu Konto magazynu
    Zasób źródłowy contosoorders
    Nazwa tematu systemowego <create any name>
    Filtruj do typów zdarzeń Utworzenie obiektu blob i usunięcie obiektu blob
  4. Wybierz Utwórz.

Testowanie subskrypcji usługi Event Grid

  1. Utwórz plik o nazwie customer-order.csv, wklej następujące informacje do tego pliku i zapisz go na komputerze lokalnym.

    InvoiceNo,StockCode,Description,Quantity,InvoiceDate,UnitPrice,CustomerID,Country
    536371,99999,EverGlow Single,228,1/1/2018 9:01,33.85,20993,Sierra Leone
    
  2. W przeglądarce Storage prześlij ten plik do folderu wejściowego swojego konta storage.

    Po przekazaniu pliku jest wywoływane zdarzenie Microsoft.Storage.BlobCreated. Usługa Event Grid powiadamia wszystkich subskrybentów o tym zdarzeniu. W tym przypadku jedynym subskrybentem jest funkcja Azure. Funkcja platformy Azure analizuje parametry zdarzenia, aby określić, które zdarzenie wystąpiło. Następnie przekazuje adres URL pliku do zadania Databricks. Zadanie Databricks odczytuje plik i dodaje wiersz do tabeli Delta Databricks, która znajduje się w Twoim koncie pamięci masowej.

  3. Aby sprawdzić, czy zadanie się powiodło, obejrzyj przebiegi dla swojego zadania. Widzisz status zakończenia. Aby uzyskać więcej informacji na temat wyświetlania uruchomień zadania, zobacz Wyświetlanie uruchomień zadania.

  4. W komórce nowego skoroszytu uruchom to zapytanie, aby zobaczyć zaktualizowaną tabelę Delta.

    %sql select * from customer_data
    

    Zwrócona tabela zawiera najnowszy rekord.

    Zrzut ekranu zapytania tabeli Databricks Delta pokazujący najnowszy rekord.

  5. Aby zaktualizować ten rekord, utwórz plik o nazwie customer-order-update.csv, wklej następujące informacje do tego pliku i zapisz go na komputerze lokalnym.

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

    Ten plik CSV jest niemal identyczny z poprzednim, z tą różnicą, że ilość zamówienia została zmieniona z 228 na .22

  6. W przeglądarce Storage prześlij ten plik do folderu wejściowego swojego konta storage.

  7. Uruchom ponownie zapytanie, select aby wyświetlić zaktualizowaną tabelę delty.

    %sql select * from customer_data
    

    Zwrócona tabela zawiera zaktualizowany rekord.

    Zrzut ekranu zapytania tabeli Databricks Delta pokazujący zaktualizowany rekord.

Czyszczenie zasobów

Gdy zasoby nie będą już potrzebne, usuń grupę zasobów i wszystkie powiązane zasoby. Aby usunąć grupę zasobów, wybierz grupę zasobów dla konta pamięci i wybierz Usuń.

Następne kroki