Tutorial: Implementar o padrão de captura do lago de dados para atualizar uma tabela Delta do Databricks

Este tutorial mostra como manipular eventos em uma conta de armazenamento que tem um namespace hierárquico.

Constroem uma pequena solução que permite preencher uma tabela Delta do Databricks carregando um ficheiro de valores separados por vírgulas (CSV) que descreve uma ordem de venda. Constrói esta solução ligando uma subscrição do Event Grid, uma Função do Azure e um Job no Azure Databricks.

Neste tutorial, irá aprender a:

  • Crie uma subscrição da Grade de Eventos que acione uma Função do Azure.
  • Crie uma Função do Azure que receba uma notificação de um evento e, em seguida, execute o trabalho no Azure Databricks.
  • Crie um trabalho Databricks que insere um pedido do cliente em uma tabela Delta do Databricks localizada na conta de armazenamento.

Constróis esta solução por ordem inversa, começando pelo espaço de trabalho do Azure Databricks.

Pré-requisitos

Criar uma ordem de venda

Primeiro, crie um ficheiro CSV que descreva uma encomenda de venda e depois carregue esse ficheiro para a conta de armazenamento. Mais tarde, usa os dados deste ficheiro para preencher a primeira linha da sua tabela Delta Databricks.

  1. Aceda à sua nova conta de armazenamento no portal do Azure.

  2. Selecione Explorador de armazenamento>Contentores de blobs>Adicionar contentor e crie um novo contentor chamado data.

    Captura de ecrã da criação de um contentor no navegador Armazenamento do Azure.

  3. No contêiner de dados, crie um diretório chamado input.

  4. Cole o texto a seguir em um editor de texto.

    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. Guarde este ficheiro no seu computador local e dê-lhe o nome de data.csv.

  6. No navegador de armazenamento, carregue este ficheiro para a pasta de entrada .

Criar um trabalho no Azure Databricks

Nesta secção, realiza estas tarefas:

  • Crie um espaço de trabalho do Azure Databricks.
  • Crie um bloco de notas.
  • Crie e preencha uma tabela Delta do Databricks.
  • Adicione código que insere linhas na tabela Delta do Databricks.
  • Crie um trabalho.

Criar uma área de trabalho do Azure Databricks

Nesta secção, cria um espaço de trabalho Azure Databricks utilizando o portal Azure.

  1. Crie um espaço de trabalho do Azure Databricks. Nomeie o espaço de trabalho contoso-orders. Consulte Criar um espaço de trabalho do Azure Databricks.

  2. Crie um cluster. Nomeie o cluster customer-order-cluster. Consulte Criar um cluster.

  3. Crie um bloco de notas. Nomeie o bloco de anotações configure-customer-table e escolha Python como o idioma padrão do bloco de anotações. Consulte Criar um bloco de notas.

Criar e preencher uma tabela Delta do Databricks

  1. No bloco de notas que criou, copie e cole o seguinte bloco de código na primeira célula, mas ainda não execute este código.

    Substitua os valores de marcador de posição appId, password e tenant neste bloco de código pelos valores que recolheu ao concluir os pré-requisitos deste tutorial.

    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'
    

    Este código cria um widget chamado source_file. Mais tarde, você criará uma Função do Azure que chama esse código e passa um caminho de arquivo para esse widget. Este código também autentica o seu principal de serviço com a conta de armazenamento e cria algumas variáveis que usa noutras células.

    Nota

    Em uma configuração de produção, considere armazenar sua chave de autenticação no Azure Databricks. Depois, adiciona uma chave de consulta ao teu bloco de código em vez da chave de autenticação.

    Por exemplo, em vez de usar esta linha de código: spark.conf.set("fs.azure.account.oauth2.client.secret", "<password>"), use a seguinte linha de código: spark.conf.set("fs.azure.account.oauth2.client.secret", dbutils.secrets.get(scope = "<scope-name>", key = "<key-name-for-service-credential>")).

    Depois de concluir este tutorial, consulte o artigo do Azure Data Lake Storage no site do Azure Databricks para ver exemplos desta abordagem.

  2. Pressione SHIFT + ENTER para executar o código neste bloco.

  3. Copie e cole o bloco de código seguinte numa célula diferente e depois pressione SHIFT + ENTER para executar o código desse bloco.

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

    Este código cria a tabela Databricks Delta na sua conta de armazenamento e depois carrega alguns dados iniciais do ficheiro CSV que carregou anteriormente.

  4. Depois de este bloco de código correr com sucesso, remova-o do seu caderno.

Adicionar código que insere linhas na tabela Delta do Databricks

  1. Copie e cole o seguinte bloco de código em uma célula diferente, mas não execute essa célula.

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

    Este código insere dados numa vista de tabela temporária utilizando dados de um ficheiro CSV. O caminho para esse ficheiro CSV vem do widget de entrada que criou numa etapa anterior.

  2. Copie e cole o bloco de código a seguir em uma célula diferente. Esse código mescla o conteúdo da exibição de tabela temporária com a tabela Delta do 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)
    

Criar um emprego

Cria um trabalho que execute o caderno que criaste anteriormente. Mais tarde, cria uma Função Azure que executa este trabalho quando um evento é gerado.

  1. Selecione Novo>Emprego.

  2. Dá um nome ao trabalho, escolhe o caderno que criaste e seleciona um cluster. Em seguida, selecione Criar para criar o trabalho.

    O novo emprego aparece na lista de Empregos com o caderno e o cluster que selecionaste.

Criar uma Função do Azure

Crie uma função Azure que execute o trabalho.

  1. No seu espaço de trabalho Azure Databricks, selecione o seu nome de utilizador Azure Databricks na barra superior. Na lista suspensa, selecione Definições do Utilizador.

  2. Na guia Tokens de acesso, selecione Gerar novo token.

  3. Copie o token que aparece e depois selecione Feito.

  4. No canto superior do espaço de trabalho Databricks, escolha o ícone de pessoas e, em seguida, escolha Configurações do usuário.

    Captura de ecrã do menu de definições do utilizador para gerar um token de acesso Databricks.

  5. Selecione o botão Gerar novo token e, em seguida, selecione o botão Gerar .

    Certifica-te de copiar o token para um local seguro. A tua Função Azure precisa que este token se autentique com o Databricks para poder executar o trabalho.

  6. A partir do menu do portal do Azure ou a partir da Home page, selecione Criar um recurso.

  7. Na página Novo, selecione Computação>Aplicação de Função.

  8. Na guia Noções básicas da página Criar aplicativo de função, escolha um grupo de recursos e altere ou verifique as seguintes configurações:

    Configuração Valor
    Nome da Aplicação de Funções contosoorder
    Pilha de tempo de execução .NET
    Publicar Código
    Sistema operativo Windows
    Tipo de plano Consumo (Sem servidor)
  9. Selecione Rever + criar e, em seguida, selecione Criar.

    Quando a implantação estiver concluída, selecione Ir para o recurso para abrir a página de visão geral do aplicativo de função.

  10. No grupo Configurações, selecione Configuração.

  11. Na página Configurações do aplicativo, escolha o botão Nova configuração do aplicativo para adicionar cada configuração.

    Captura de ecrã da adição de uma nova definição de aplicação na configuração da Function App.

    Adicione as seguintes configurações:

    Nome da configuração Valor
    DBX_INSTANCE A região do seu espaço de trabalho databricks. Por exemplo: westus2.azuredatabricks.net
    DBX_PAT O token de acesso pessoal que você gerou anteriormente.
    DBX_JOB_ID O identificador do trabalho em execução.
  12. Selecione Salvar para confirmar essas configurações.

  13. No grupo Funções, selecione Funções e, em seguida, selecione Criar.

  14. Escolha Azure Event Grid Trigger.

    Instale a extensão Microsoft.Azure.WebJobs.Extensions.EventGrid se for solicitado. Se precisares de instalar, seleciona novamente o Azure Event Grid Trigger para criar a função.

    O painel Nova função é exibido.

  15. Em Nova Função, introduza UpsertOrder o nome da função e depois selecione Criar.

  16. Substitua o conteúdo do ficheiro de código pelo seguinte código e depois selecione Guardar:

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

    Este código analisa informações sobre o evento de armazenamento que foi gerado e depois cria uma mensagem de pedido com a URL do ficheiro que desencadeou o evento. Como parte da mensagem, a função passa um valor para o widget source_file que você criou anteriormente. O código da função envia a mensagem para o trabalho Databricks e usa o token que obtiveste anteriormente como autenticação.

Criar uma subscrição do Event Grid

Nesta secção, cria uma subscrição de Event Grid que invoca a Função do Azure quando os ficheiros são carregados para a conta de armazenamento.

  1. Selecione Integração. Na página Integração, selecione Acionador do Event Grid.

  2. No painel Editar gatilho, nomeie o evento eventGridEvente selecione Criar assinatura de evento.

    Nota

    O nome eventGridEvent corresponde ao nome do parâmetro que a Função Azure recebe.

  3. Na guia Noções básicas da página Criar Assinatura de Evento, altere ou verifique as seguintes configurações:

    Configuração Valor
    Nome subscrição-de-eventos-de-ordens-contoso
    Tipo de tópico Conta de armazenamento
    Recurso de origem contosoorders
    Nome do tópico do sistema <create any name>
    Filtrar Tipos de Eventos Blob criado e Blob excluído
  4. Selecione Criar.

Testar a subscrição da Grelha de Eventos

  1. Crie um ficheiro com o nome customer-order.csv, cole as seguintes informações nesse ficheiro e guarde-o no computador 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. No navegador de armazenamento, carregue este ficheiro na pasta de entrada da sua conta de armazenamento.

    Quando carrega um ficheiro, despoleta o evento Microsoft.Storage.BlobCreated. A Grelha de Eventos notifica todos os subscritores desse evento. Neste caso, a Função Azure é o único assinante. A Função do Azure analisa os parâmetros de evento para determinar qual evento ocorreu. Depois, passa o URL do ficheiro para a tarefa do Databricks. A tarefa do Databricks lê o ficheiro e adiciona uma linha à tabela Delta do Databricks localizada na sua conta de armazenamento.

  3. Para verificar se o trabalho foi bem-sucedido, consulte as execuções do seu trabalho. Vê o estado de conclusão. Para obter mais informações sobre como ver as execuções de uma tarefa, consulte Ver as execuções de uma tarefa.

  4. Numa nova célula do livro de exercícios, execute esta consulta para ver a tabela delta atualizada.

    %sql select * from customer_data
    

    A tabela retornada mostra o registro mais recente.

    Captura de ecrã da consulta da tabela Delta do Databricks que mostra o registo mais recente.

  5. Para atualizar esse registro, crie um arquivo chamado customer-order-update.csv, cole as seguintes informações nesse arquivo e salve-o no computador local.

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

    Este ficheiro CSV é quase idêntico ao anterior, exceto que a quantidade da ordem é alterada de 228 para 22.

  6. No navegador de armazenamento, carregue este ficheiro na pasta de entrada da sua conta de armazenamento.

  7. Execute a select consulta novamente para ver a tabela delta atualizada.

    %sql select * from customer_data
    

    A tabela retornada mostra o registro atualizado.

    Captura de ecrã da consulta da tabela Delta do Databricks mostrando o registo atualizado.

Limpar recursos

Quando já não precisares dos recursos, elimina o grupo de recursos e todos os recursos relacionados. Para eliminar o grupo de recursos, selecione o grupo de recursos para a conta de armazenamento e selecione Eliminar.

Próximos passos