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.
- Åpne Azure Portal, og søk etter «Azure Databricks».

- Klikk på Create Azure Databricks service.
- Velg subscription og resource group der de andre ressursene dine, for eksempel Azure Data Factory, befinner seg.
- Velg et beskrivende workspace-navn.
- Velg regionen nærmest deg.
- Velg Trial som pricing tier, med mindre dette er en produksjonsarbeidsbelastning. Velg i så fall Premium.
- Gå til Review + create.
- Create
Opprette service principal slik at Databricks får tilgang til Azure Key Vault
- Åpne en PowerShell Core-terminal i Windows Terminal eller PowerShell Core som administrator, eller med utvidede tillatelser.
- 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
- Importer Az PowerShell-modulen i den gjeldende økten.
- 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.
- 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>'
- Opprett Azure AD Service Principal.
##Opprett Azure AD Service Principal
$createServicePrincipal = New-AzADServicePrincipal -DisplayName $servicePrincipalName -Role 'Storage Blob Data Contributor' -Scope $servicePrincipalScope
- 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.)
- Utvid menyen Workspace i venstre panel, og åpne deretter mappen Shared. Klikk på nedoverpilen i den delte mappen, og velg Create Notebook.

- 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. 
- 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")
- 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"
- 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)
- 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)
- 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)
- 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)
- 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"))
- 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)
- 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)
- 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.
- Åpne Azure Key Vault i portalen, og gå til bladet Access Policies.
- Klikk på +Add Access Policy.

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

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

- 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.
- Åpne en ny fane, og gå til Databricks-workspacet ditt. URL-en skal ligne på
https://adb-<workspaceid>.<randomid>.azuredatabricks.net/. - 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.
- Gå til Key Vault-en din i Azure Portal i en annen fane, og åpne bladet Properties.
- Noter følgende verdier:
- Gå tilbake til Databricks-workspace-fanen med siden for å opprette secret scope.
- 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.
- I eksempelmiljøet mitt kreves ikke høyeste sikkerhetsnivå, så jeg tillater alle brukere å administrere principal-en i rullegardinlisten nedenfor.
- Lim inn Vault URI i feltet DNS Name.
- Lim inn Resource ID i feltet Resource ID.
- Klikk på Create.
- 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
- Gå til Azure Databricks-workspacet i Azure Portal i en ny fane.
- Åpne bladet Access control (IAM).
- Klikk på Add.
- Klikk på Add role assignment.

- 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.
- Klikk på Next.
- På fanen Members setter du alternativknappen ved siden av Assign access to til Managed identity.
- Klikk på Select members.
- 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).

- Velg riktig data factory, og klikk deretter på Select.
- Klikk på Review + assign og deretter Assign for å fullføre prosessen.
Legge til Azure Databricks som en Linked Service i Azure Data Factory
- Åpne Azure Data Factory studio som du har brukt i denne øvelsen, i en ny fane.
- Klikk på ikonet
. - Klikk på bladet Linked services.
- Klikk på New.
- Klikk på fanen Compute, velg Azure Databricks, og klikk deretter på Continue.

- Gi linked service navnet LS_AzureDatabricks, eller et annet navn etter navnekonvensjonen din.
- 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.

- Klikk på Test connection for å kontrollere at den fungerer.
- Klikk på Save.
- Klikk på Publish all.

- Du skal nå se LS_AzureDatabricks i listen over linked services.

Legge Azure Databricks-notebooken til i den eksisterende Azure Data Factory-pipelinjen
- Gå til author mode i Azure Data Factory studio ved å klikke på ikonet
. - Finn den eksisterende pipelinen som ble opprettet tidligere, og åpne den.
- Utvid alle aktivitetene.

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

- Klikk på fanen Azure Databricks for å knytte aktiviteten til den nylig opprettede linked service-en.

- Klikk på fanen Settings, og klikk deretter på Browse.
- Søk etter notebooken du lagret tidligere, og klikk på OK.

- Til slutt må vi angi base parameters for Databricks-aktiviteten, slik at folderPath og pipelineRunId kan sendes til notebooken ved kjøring.
- Utvid Base parameters.
- Klikk på +New to ganger.
- Gi den første parameteren navnet filename.
- Angi en dynamic content-verdi med følgende kode:
- Gi den andre parameteren navnet pipelineRunId.
- Angi en dynamic content-verdi med følgende kode:
- Klikk på knappen Publish all.
- 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.