Del 3: Rapportering på Log Analytics-data | Behandle data med Azure Databricks

| 12 min lesing

I dette neste innlegget i serien transformeres dataene som er hentet ut med Azure Data Factory, til et lesbart format med Azure Databricks. Deretter legges denne logikken inn i Azure Data Factory-pipelinjen.

Hvis du gikk glipp av det, finner du en lenke til del 2 her.

Opprette Azure Databricks-workspace

For å fullføre neste del av prosessen må et Azure Databricks-workspace distribueres i abonnementet som brukes.

  1. Åpne Azure Portal, og søk etter «Azure Databricks».
  2. Klikk på Create Azure Databricks service.
  3. Velg subscription og resource group der de andre ressursene dine, for eksempel Azure Data Factory, befinner seg.
  4. Velg et beskrivende workspace-navn.
  5. Velg regionen nærmest deg.
  6. Velg Trial som pricing tier, med mindre dette er en produksjonsarbeidsbelastning. Velg i så fall Premium.
  7. Gå til Review + create.
  8. Create

Opprette service principal slik at Databricks får tilgang til Azure Key Vault

  1. Åpne en PowerShell Core-terminal i Windows Terminal eller PowerShell Core som administrator, eller med utvidede tillatelser.
  2. Installer Az PowerShell-modulen. Hvis den allerede er installert, kontrollerer du om det finnes oppdateringer. Trykk A og Enter for å stole på repositoriet mens de nyeste modulene lastes ned og installeres. Dette kan ta noen minutter.
    ##Kjør hvis Az PowerShell-modulen ikke allerede er installert
     Install-Module Az -AllowClobber
    
    ##Kjør hvis Az PowerShell-modulen allerede er installert, men oppdateringskontroll er nødvendig
    Update-Module Az
    
  3. Importer Az PowerShell-modulen i den gjeldende økten.
    Import-Module Az
    
  4. Autentiser mot Azure som vist nedenfor. Hvis tenant eller subscription er en annen enn standarden som er knyttet til Azure AD-identiteten din, kan du angi dette på samme kodelinje med parameterne -Tenant eller -Subscription. For å finne Tenant GUID åpner du Azure Portal, kontrollerer at du er logget på riktig tenant og går til Azure Active Directory. Directory GUID skal vises på startsiden. For å finne Subscription GUID bytter du til riktig tenant i Azure Portal, skriver Subscriptions i søkefeltet og trykker Enter. GUID-ene for alle subscriptions du har tilgang til i den tenant-en, vises.
    Connect-AzAccount
    
  5. Angi først noen variabler som skal brukes i de følgende kommandoene.
    ##Angi variablene våre
    $vaultName = '<skriv inn navnet på tidligere opprettet Key Vault>'
    $databricksWorkspaceName = '<skriv inn navnet på workspacet ditt>'
    $servicePrincipalName = $databricksWorkspaceName, '-ServicePrincipal'
    $servicePrincipalScope = '/subscriptions/<subscriptionid>/resourceGroups/<resourcegroupname>/<providers/Microsoft.Storage/storageAccounts/<storageaccountname>'
    
  6. Opprett Azure AD Service Principal.
    ##Opprett Azure AD Service Principal
    $createServicePrincipal = New-AzADServicePrincipal -DisplayName $servicePrincipalName -Role 'Storage Blob Data Contributor' -Scope $servicePrincipalScope
    
  7. Lagre til slutt client ID og client secret i Key Vault.
    ##Lagre service principal client ID i Key Vault
    Set-AzKeyVaultSecret -VaultName $vaultName -Name 'databricks-service-principal-client-id' -SecretValue (ConvertTo-SecureString -AsPlainText ($createServicePrincipal.ApplicationId))
    
    ##Lagre service principal secret i Key Vault
    Set-AzKeyVaultSecret -VaultName $vaultName -Name 'databricks-service-principal-secret' -SecretValue $createServicePrincipal.Secret
    

Opprette notebooken

NB: Det ble innført en brytende endring i Spark 3.2 som gjør at koden i steg 6 ikke fungerer. Bruk Databricks 9.1 LTS Runtime.

(All koden som vises, ligger i GitHub-repositoriet mitt – her.)

  1. Utvid menyen Workspace i venstre panel, og åpne deretter mappen Shared. Klikk på nedoverpilen i den delte mappen, og velg Create Notebook.
  2. Velg et navn for notebooken som vist nedenfor. Jeg bruker loganalytics-fileprocessor. Scala er språket som brukes i denne notebooken, så velg det i rullegardinlisten Default Language, og klikk på Create.
  3. NB: Vær forsiktig når du kopierer og limer inn kode. Ekstra eller inkompatibel mellomromstegn kan bli limt inn i Databricks-notebooken og føre til feil under kjøring. Nå som notebooken er opprettet, kan vi skrive den første kodeblokken for å gjøre den dynamisk. Dette gjør to ting: pipelineRunId fra Data Factory og mappebanen der filen er skrevet, kan sendes inn som parametere. Gi kodeblokken tittelen «Parameters» med rullegardinmenyen på høyre side av kodeblokken.
    //Definer folderPath
    dbutils.widgets.text("rawFilename", "")
    val paramRawFile = dbutils.widgets.get("rawFilename")
    
    //Definer pipelineRunId
    dbutils.widgets.text("pipelineRunId", "")
    val paramPipelineRunId = dbutils.widgets.get("pipelineRunId")
    
  4. Opprett en ny kodeblokk kalt «Declare and set variables». Deretter deklarerer vi variablene som skal brukes i notebooken. Legg merke til at jeg bruker abfss-driveren, som er optimalisert for Data Lake Gen 2 og big data-arbeidsbelastninger, i stedet for standard wasbs for vanlig Blob Storage. I dette eksemplet er datasettet ikke spesielt stort, men ytelsesforskjellen vil merkes når datasettet blir betydelig større, så det er best å følge anbefalt praksis.
    val dataLake = "abfss://<containername>@<datalakename>.dfs.core.windows.net"
    val rawFolderPath = ("/raw-data/")
    val rawFullPath = (dataLake + rawFolderPath + paramRawFile)
    val outputFolderPath = "/output-data/"
    val databricksServicePrincipalClientId = dbutils.secrets.get(scope = "databricks-secret-scope", key = "databricks-service-principal-client-id")
    val databricksServicePrincipalClientSecret = dbutils.secrets.get(scope = "databricks-secret-scope", key = "databricks-service-principal-secret")
    val azureADTenant = dbutils.secrets.get(scope = "databricks-secret-scope", key = "azure-ad-tenant-id")
    val endpoint = "https://login.microsoftonline.com/" + azureADTenant + "/oauth2/token"
    val dateTimeFormat = "yyyy_MM_dd_HH_mm"
    
  5. Opprett en ny kodeblokk kalt «Set storage context and read source data». For å få tilgang til data lake-en må vi angi storage context for økten når notebooken kjøres.
    import org.apache.spark.sql
    
    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", databricksServicePrincipalClientId)
    spark.conf.set("fs.azure.account.oauth2.client.secret", databricksServicePrincipalClientSecret)
    spark.conf.set("fs.azure.account.oauth2.client.endpoint", endpoint)
    
    val sourceDf = spark.read.option("multiline",true).json(rawFullPath)
    
  6. Opprett en ny kodeblokk kalt «Explode source data columns into tabular format». Nå kan vi gjøre dataene om til tabellformat med en explode-funksjon. Erstatt kolonnene som er eksplisitt definert på linje 9, slik at de samsvarer med output-skjemaet fra Log Analytics-function-en du har valgt.
    import org.apache.spark.sql.functions._
    
    val explodedDf = sourceDf.select(explode($"tables").as("tables"))
    .select($"tables.columns".as("column"), explode($"tables.rows").as("row"))
    .selectExpr("inline(arrays_zip(column, row))")
    .groupBy()
    .pivot($"column.name")
    .agg(collect_list($"row"))
    .selectExpr("inline(arrays_zip(timeGenerated, userAction, appUrl, successFlag, httpResultCode, durationOfRequestMs, clientType, clientOS, clientCity, clientStateOrProvince, clientCountryOrRegion, clientBrowser, appRoleName, snapshotTimestamp))")
    
    display(explodedDf)
    
  7. Opprett en ny kodeblokk kalt «Transform source data». Vi må legge til noen datapunkter i data frame-en for å kunne partisjonere effektivt. Når datasettet vokser etter hvert som flere pipeline-kjøringer skjer, minimerer dette ytelsesflaskehalser der det er mulig. I dette tilfellet legger vi også til ADF pipeline run ID for å gjøre feilsøking enklere dersom det oppstår datakvalitetsproblemer i målet.
    import org.apache.spark.sql.SparkSession
    
    val pipelineRunIdSparkSession = SparkSession
        .builder()
        .appName("Pipeline Run Id Appender")
        .getOrCreate()
    
    // Registrer den transformerte DataFrame-en som en midlertidig SQL-visning
    
    explodedDf.createOrReplaceTempView("transformedDf")
    
    val transformedDf = spark.sql("""
    SELECT DISTINCT
    timeGenerated,
    userAction,
    appUrl,
    successFlag,
    httpResultCode,
    durationOfRequestMs,
    clientType,
    clientOS,
    clientCity,
    clientStateOrProvince,
    clientCountryOrRegion,
    clientBrowser,
    appRoleName,
    snapshotTimestamp,
    '""" + paramPipelineRunId + """' AS pipelineRunId,
    CAST(timeGenerated AS DATE) AS requestDate,
    HOUR(timeGenerated) AS requestHour
    FROM transformedDf""").toDF
    
    display(transformedDf)
    
  8. Opprett en ny kodeblokk kalt «Write transformed data to delta lake». For å beholde historikk i faktadataene kan vi bruke Delta Lake-funksjonaliteten i Azure Databricks. Først må vi kontrollere at en database med navnet logAnalyticsdb finnes. Hvis ikke, oppretter den første kodelinjen i neste blokk den. Selv om det ikke er obligatorisk å bruke eksplisitt navngitte databaser i Databricks, blir det mye enklere når mange tabeller lagres. Dette hindrer at default-databasen blir uoversiktlig. Tilnærmingen ligner organiseringen av en Microsoft SQL Server-database, der du kategoriserer tabeller i forskjellige skjemaer, for eksempel raw, staging og final, i stedet for å ha alt i dbo-skjemaet.
    import org.apache.spark.sql.SaveMode
    
    display(spark.sql("CREATE DATABASE IF NOT EXISTS logAnalyticsdb"))
    
    transformedDf.write
      .format("delta")
      .mode("append")
      .option("mergeSchema","true")
      .partitionBy("requestDate", "requestHour")
      .option("path", "/delta/logAnalytics/websiteLogs")
      .saveAsTable("logAnalyticsdb.websitelogs")
    
    val transformedDfDelta = spark.read.format("delta")
      .load("/delta/logAnalytics/websiteLogs")
    
    display(transformedDfDelta)
    
  9. Opprett en ny kodeblokk kalt «Remove stale data, optimize and vacuum delta table». For å kontrollere størrelsen på Delta Lake-tabellen må vi gjøre følgende: Først sletter vi poster som er eldre enn 7 dager. Deretter optimaliserer vi tabellen ved å utføre en ZORDER-funksjon på kolonnen snapshotTimestamp og samle dataene i så få parquet-filer som mulig, for å forbedre spørringsytelsen. Til slutt kjører vi vacuum for å sikre at bare nødvendige parquet-filer beholdes. Her ønsker vi å beholde data for 7 dager. Syntaksen for vacuum-kommandoen forventer parameteren angitt i timer. I vårt tilfelle passer standardinnstillingen på 7 dager / 168 timer, så vi trenger ingen parameter. Deretter viser vi Delta-historikken.
    import io.delta.tables._
    
    display(spark.sql("""
    DELETE
    FROM logAnalyticsdb.websitelogs
    WHERE requestDate < DATE_ADD(CURRENT_TIMESTAMP, -7)
    """))
    
    display(spark.sql("""
    OPTIMIZE logAnalyticsdb.websitelogs
    ZORDER BY (snapshotTimestamp)
    """))
    
    val deltaTable = DeltaTable.forPath(spark, "/delta/logAnalytics/websiteLogs")
    deltaTable.vacuum()
    
    display(spark.sql("DESCRIBE HISTORY logAnalyticsdb.websitelogs"))
    
  10. Opprett en ny kodeblokk kalt «Create data frame for output file». Velg nå datapunktene fra Delta Lake-tabellen som vi trenger i output-filen.
    import org.apache.spark.sql.SparkSession
    
    val outputFileSparkSession = SparkSession
    .builder()
    .appName("Output File Generator")
    .getOrCreate()
    
    val outputDf = spark.sql("""
    SELECT DISTINCT
    timeGenerated,
    userAction,
    appUrl,
    successFlag,
    httpResultCode,
    durationOfRequestMs,
    clientType,
    clientOS,
    clientCity,
    clientStateOrProvince,
    clientCountryOrRegion,
    clientBrowser,
    appRoleName,
    requestDate,
    requestHour,
    pipelineRunId
    FROM logAnalyticsdb.websitelogs
    WHERE pipelineRunId = '""" + paramPipelineRunId + """'
    """
    )
    display(outputDf)
    
  11. Opprett en ny kodeblokk kalt «Create output file in data lake». La oss lagre denne output-en i data lake-en.
    import org.apache.spark.sql.functions._
    
    val currentDateTimeLong = current_timestamp().expr.eval().toString.toLong
    val currentDateTime = new java.sql.Timestamp(currentDateTimeLong/1000).toLocalDateTime.format(java.time.format.DateTimeFormatter.ofPattern(dateTimeFormat))
    val outputSessionFolderPath = ("appLogs_" + currentDateTime)
    val fullOutputPath = (dataLake + outputFolderPath + outputSessionFolderPath)
    
    outputDf.write.parquet(fullOutputPath)
    
  12. Opprett en ny kodeblokk kalt «Display output parameters». Vi må legge til én kodeblokk til, som blir en output-parameter vi sender tilbake til Azure Data Factory, slik at den vet hvilken fil som skal behandles.
    dbutils.notebook.exit(outputFolderPath + outputSessionFolderPath)
    

Gi Azure Databricks tillatelser i Azure Key Vault

For at Azure Databricks skal kunne lese secrets som er lagret i Azure Key Vault, må eksplisitte tillatelser angis for Azure Databricks ved hjelp av en Access Policy.

NB: Azure role-based access control er den vanlige metoden for å tildele tillatelser i Azure, men her ønsker vi svært spesifikke tillatelser. En RBAC-rolle for Key Vault kan være for bredt definert.

  1. Åpne Azure Key Vault i portalen, og gå til bladet Access Policies.
  2. Klikk på +Add Access Policy.
  3. Angi nødvendig tillatelsesnivå. I dette tilfellet trenger Azure Databricks bare å kunne Get og List secrets som er lagret i Key Vault – ikke mer og ikke mindre.
  4. Til slutt må vi angi principal, altså identiteten til parten som får disse tillatelsene. Da dette ble skrevet, fikk ikke hvert enkelt Databricks-workspace sin egen Azure Active Directory service principal. Du må derfor bruke principal-en for den globale Azure Databricks Enterprise Application i Azure Active Directory. Klikk på None Selected ved siden av Select principal for å velge den. Søk etter AzureDatabricks; det skal bare være én oppføring. Velg elementet i listen, klikk på Select, og klikk deretter på Add.
  5. Du kommer da tilbake til bladet Access Policies. Klikk på Save for å bruke endringene.

Opprette secret scope i Azure Databricks

Nå som vi har opprettet tillatelser på Key Vault for den globale AzureDatabricks enterprise application, må vi konfigurere workspacet til å koble til Key Vault og hente secrets ved å opprette et secret scope.

NB: Du trenger minst Contributor-tillatelser på selve Key Vault-en for å utføre denne operasjonen.

  1. Åpne en ny fane, og gå til Databricks-workspacet ditt. URL-en skal ligne på https://adb-<workspaceid>.<randomid>.azuredatabricks.net/.
  2. Legg til #secrets/createScope i URL-en, slik at den blir https://adb-<workspaceid>.<randomid>.azuredatabricks.net/#secrets/createScope.
    • Merk at denne ekstra delen av URL-en skiller mellom store og små bokstaver.
  3. Gå til Key Vault-en din i Azure Portal i en annen fane, og åpne bladet Properties.
  4. Noter følgende verdier:
    • Vault URI
    • Resource ID
  5. Gå tilbake til Databricks-workspace-fanen med siden for å opprette secret scope.
  6. Angi scope-navnet som verdien du definerte på linje 5 i del 4 av opprettelsen av notebooken. I eksemplet bruker jeg databricks-secret-scope. I mer realistiske scenarioer med ulike miljøer, som Development, Test, Production og Quality Assurance, vil du sannsynligvis ha Key Vault-er avgrenset til hvert miljø med et tilsvarende secret scope.
  7. I eksempelmiljøet mitt kreves ikke høyeste sikkerhetsnivå, så jeg tillater alle brukere å administrere principal-en i rullegardinlisten nedenfor.
  8. Lim inn Vault URI i feltet DNS Name.
  9. Lim inn Resource ID i feltet Resource ID.
  10. Klikk på Create.
  11. Hvis brukeren som utfører operasjonen har minst Contributor-tillatelser, skal følgende melding vises.

Koble Azure Databricks til Azure Data Factory

Nå som Databricks-miljøet er konfigurert og klart, må vi først gi Azure Data Factory tilgang til Databricks-workspacet og deretter integrere notebooken i Azure Data Factory-pipelinen.

Gi Azure Data Factory tillatelser på Azure Databricks-workspaceet

  1. Gå til Azure Databricks-workspacet i Azure Portal i en ny fane.
  2. Åpne bladet Access control (IAM).
  3. Klikk på Add.
  4. Klikk på Add role assignment.
  5. Klikk på raden med navnet Contributor. Azure Databricks har dessverre ikke et eget sett med spesifikke roller for Azure Databricks, så den generiske Contributor-rollen må brukes.
  6. Klikk på Next.
  7. På fanen Members setter du alternativknappen ved siden av Assign access to til Managed identity.
  8. Klikk på Select members.
  9. Finn Azure Data Factory-instansen din ved å velge riktig subscription og deretter filtrere på managed identity-type. I dette tilfellet er det Data factory (V2).
  10. Velg riktig data factory, og klikk deretter på Select.
  11. Klikk på Review + assign og deretter Assign for å fullføre prosessen.

Legge til Azure Databricks som en Linked Service i Azure Data Factory

  1. Åpne Azure Data Factory studio som du har brukt i denne øvelsen, i en ny fane.
  2. Klikk på ikonet Azure Data Factory Manage-ikon, verktøykasse med skiftenøkkel.
  3. Klikk på bladet Linked services.
  4. Klikk på New.
  5. Klikk på fanen Compute, velg Azure Databricks, og klikk deretter på Continue. Azure Data Factory viser menyen for å opprette en ny linked compute service, med Azure Databricks valgt
  6. Gi linked service navnet LS_AzureDatabricks, eller et annet navn etter navnekonvensjonen din.
  7. Konfigurer den slik:
    • Bruk standard Integration runtime (AutoResolveIntegrationRuntime), med mindre du har et bestemt behov for en annen.
    • Velg From Azure subscription i rullegardinlisten for account selection method.
    • Velg riktig Azure subscription, og deretter riktig workspace.
    • Velg en ny job cluster som cluster type for å redusere kostnadene.
    • Angi autentiseringstype til Managed service identity.
    • Workspace resource ID skal være forhåndsutfylt.
    • Velg cluster version 9.1 LTS. Det ble innført en brytende endring i Spark 3.2 som gjør at koden i steg 6 ikke fungerer. Bruk Databricks 9.1 LTS Runtime.
    • Standard_DS3_v2 skal passe som cluster node type.
    • Velg Python version 3.
    • Angi autoskalering av workers til minimum 1 og maksimum 2. Azure Data Factory viser forventet konfigurasjon for Azure Databricks linked service, del 1 Azure Data Factory viser forventet konfigurasjon for Azure Databricks linked service, del 2
  8. Klikk på Test connection for å kontrollere at den fungerer.
  9. Klikk på Save.
  10. Klikk på Publish all. Azure Data Factory-knappen Publish all
  11. Du skal nå se LS_AzureDatabricks i listen over linked services. Azure Data Factory linked services med LS_AzureDatabricks lagt til

Legge Azure Databricks-notebooken til i den eksisterende Azure Data Factory-pipelinjen

  1. Gå til author mode i Azure Data Factory studio ved å klikke på ikonet Azure Data Factory Author-ikon, blyant.
  2. Finn den eksisterende pipelinen som ble opprettet tidligere, og åpne den.
  3. Utvid alle aktivitetene. Eksempel på Azure Data Factory authoring view med eksempelpipelinen og tilgjengelige aktiviteter
  4. Dra en Databricks-aktivitet av typen Notebook inn i pipelinen, og legg den til som etterfølgende aktivitet etter Copy Raw Data. Kall aktiviteten Transform Source Data. Eksempel på Azure Data Factory-pipeline med den nylig tilføyde Azure Databricks notebook-aktiviteten Transform Source Data
  5. Klikk på fanen Azure Databricks for å knytte aktiviteten til den nylig opprettede linked service-en. Innstillinger for Azure Data Factory Databricks notebook-aktivitet som knytter den til Azure Databricks linked service
  6. Klikk på fanen Settings, og klikk deretter på Browse.
  7. Søk etter notebooken du lagret tidligere, og klikk på OK. Azure Data Factory Databricks notebook-aktivitet med den tidligere opprettede notebooken valgt for notebook path-innstillingen
  8. Til slutt må vi angi base parameters for Databricks-aktiviteten, slik at folderPath og pipelineRunId kan sendes til notebooken ved kjøring.
  9. Utvid Base parameters.
  10. Klikk på +New to ganger.
  11. Gi den første parameteren navnet filename.
  12. Angi en dynamic content-verdi med følgende kode:
    @variables('filename')
    
  13. Gi den andre parameteren navnet pipelineRunId.
  14. Angi en dynamic content-verdi med følgende kode:
    @pipeline().RunId
    
  15. Klikk på knappen Publish all.
  16. Kjør Debug for pipelinen for å kontrollere at alt fungerer. Merk at hvis kildedatasettet er stort, kan Databricks-notebooken bruke noen minutter på behandling etter at clusteret er opprettet. Jeg opplevde en flaskehals i steget der dataene lagres i Delta Lake. Det kan være verdt å øke antallet workers som er tilgjengelige for job clusteret, men da må Azure subscription ha riktig core-kvote for VM Series som brukes av Databricks-clusteret. De fleste VM Series har som standard 10 cores. Du kan be om kvoteøkning ved å sende inn en support request.

Oppsummering

Da er vi i mål: Vi har laget Databricks-notebooken som behandler Log Analytics-dataene, og integrert notebooken i den eksisterende Azure Data Factory-pipelinjen.