I dette neste innlegget i serien blir dataa som er henta ut med Azure Data Factory, transformerte til eit lesbart format med Azure Databricks. Deretter blir denne logikken lagd inn i Azure Data Factory-pipelinen.
Viss du gjekk glipp av det, finn du ei lenkje til del 2 her.
Opprette Azure Databricks-workspace
For å fullføre neste del av prosessen må eit Azure Databricks-workspace distribuerast i abonnementet som blir brukt.
- Opne Azure Portal, og søk etter «Azure Databricks».

- Klikk på Create Azure Databricks service.
- Vel subscription og resource group der dei andre ressursane dine, til dømes Azure Data Factory, ligg.
- Vel eit beskrivande workspace-namn.
- Vel regionen nærast deg.
- Vel Trial som pricing tier, med mindre dette er ei produksjonsarbeidsbelastning. Vel i så fall Premium.
- Gå til Review + create.
- Create
Opprette service principal slik at Databricks får tilgang til Azure Key Vault
- Opne ein PowerShell Core-terminal i Windows Terminal eller PowerShell Core som administrator, eller med utvida løyve.
- Installer Az PowerShell-modulen. Viss han alt er installert, kontrollerer du om det finst oppdateringar. Trykk A og Enter for å stole på repositoriet medan dei nyaste modulane blir lasta ned og installerte. Dette kan ta nokre minutt.
##Køyr viss Az PowerShell-modulen ikkje alt er installert
Install-Module Az -AllowClobber
##Køyr viss Az PowerShell-modulen alt er installert, men oppdateringskontroll er nødvendig
Update-Module Az
- Importer Az PowerShell-modulen i den gjeldande økta.
- Autentiser mot Azure som vist nedanfor. Viss tenant eller subscription er ein annan enn standarden som er knytt til Azure AD-identiteten din, kan du angi dette på same kodelinje med parametrane -Tenant eller -Subscription. For å finne Tenant GUID opnar du Azure Portal, kontrollerer at du er pålogga rett tenant og går til Azure Active Directory. Directory GUID skal visast på startsida. For å finne Subscription GUID byter du til rett tenant i Azure Portal, skriv Subscriptions i søkjefeltet og trykkjer Enter. GUID-ane for alle subscriptions du har tilgang til i den tenant-en, blir viste.
- Angi først nokre variablar som skal brukast i dei følgjande kommandoane.
##Angi variablane våre
$vaultName = '<skriv inn namnet på tidlegare oppretta Key Vault>'
$databricksWorkspaceName = '<skriv inn namnet 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 blei innført ei brytande endring i Spark 3.2 som gjer at koden i steg 6 ikkje fungerer. Bruk Databricks 9.1 LTS Runtime.
(All koden som blir vist, ligg i GitHub-repositoriet mitt – her.)
- Utvid menyen Workspace i venstre panel, og opne deretter mappa Shared. Klikk på nedoverpila i den delte mappa, og vel Create Notebook.

- Vel eit namn for notebooken som vist nedanfor. Eg bruker loganalytics-fileprocessor. Scala er språket som blir brukt i denne notebooken, så vel det i rullegardinlista Default Language, og klikk på Create.

- NB: Ver varsam når du kopierer og limer inn kode. Ekstra eller inkompatible mellomromsteikn kan bli limte inn i Databricks-notebooken og føre til feil under køyring. No som notebooken er oppretta, kan vi skrive den første kodeblokka for å gjere ho dynamisk. Dette gjer to ting: pipelineRunId frå Data Factory og mappebana der fila er skriven, kan sendast inn som parametrar. Gi kodeblokka tittelen «Parameters» med rullegardinmenyen på høgre side av kodeblokka.
//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 ei ny kodeblokk kalla «Declare and set variables». Deretter deklarerer vi variablane som skal brukast i notebooken. Legg merke til at eg bruker abfss-driveren, som er optimalisert for Data Lake Gen 2 og big data-arbeidsbelastningar, i staden for standard wasbs for vanleg Blob Storage. I dette eksempelet er datasettet ikkje særleg stort, men ytingsforskjellen vil merkast når datasettet blir vesentleg større, så det er best å følgje 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 ei ny kodeblokk kalla «Set storage context and read source data». For å få tilgang til data lake-en må vi angi storage context for økta når notebooken blir køyrd.
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 ei ny kodeblokk kalla «Explode source data columns into tabular format». No kan vi gjere dataa om til tabellformat med ein explode-funksjon. Erstatt kolonnane som er eksplisitt definerte på linje 9, slik at dei samsvarer med output-skjemaet frå Log Analytics-function-en du har valt.
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 ei ny kodeblokk kalla «Transform source data». Vi må leggje til nokre datapunkt i data frame-en for å kunne partisjonere effektivt. Når datasettet veks etter kvart som fleire pipeline-køyringar skjer, minimerer dette ytingsflaskehalsar der det er mogleg. I dette tilfellet legg vi òg til ADF pipeline run ID for å gjere feilsøking enklare dersom det oppstår datakvalitetsproblem 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 ei mellombels SQL-vising
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 ei ny kodeblokk kalla «Write transformed data to delta lake». For å behalde historikk i faktadataa kan vi bruke Delta Lake-funksjonaliteten i Azure Databricks. Først må vi kontrollere at ein database med namnet logAnalyticsdb finst. Viss ikkje, opprettar den første kodelinja i neste blokk han. Sjølv om det ikkje er obligatorisk å bruke eksplisitt namngjevne databasar i Databricks, blir det mykje enklare når mange tabellar blir lagra. Dette hindrar at default-databasen blir uoversiktleg. Tilnærminga liknar organiseringa av ein Microsoft SQL Server-database, der du kategoriserer tabellar i ulike skjema, til dømes raw, staging og final, i staden 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 ei ny kodeblokk kalla «Remove stale data, optimize and vacuum delta table». For å kontrollere storleiken på Delta Lake-tabellen må vi først slette postar eldre enn 7 dagar. Deretter optimaliserer vi tabellen med ein ZORDER-funksjon på kolonnen snapshotTimestamp og samlar dataa i så få parquet-filer som mogleg for å betre spørringsytinga. Til slutt køyrer vi vacuum for å sikre at berre nødvendige parquet-filer blir behaldne. Her ønskjer vi å behalde data for 7 dagar. Vacuum-kommandoen forventar parameteren i timar, og standardinnstillinga på 7 dagar / 168 timar passar i dette tilfellet. 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 ei ny kodeblokk kalla «Create data frame for output file». Vel no datapunkta frå Delta Lake-tabellen som vi treng i output-fila.
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 ei ny kodeblokk kalla «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 ei ny kodeblokk kalla «Display output parameters». Vi må leggje til éi kodeblokk til, som blir ein output-parameter vi sender tilbake til Azure Data Factory, slik at han veit kva fil som skal behandlast.
dbutils.notebook.exit(outputFolderPath + outputSessionFolderPath)
Gje Azure Databricks løyve i Azure Key Vault
For at Azure Databricks skal kunne lese secrets som er lagra i Azure Key Vault, må eksplisitte løyve spesifiserast for Azure Databricks med ein Access Policy.
NB: Azure role-based access control er den vanlege metoden for å tildele løyve i Azure, men her ønskjer vi svært spesifikke løyve. Ei RBAC-rolle for Key Vault kan vere for breitt definert.
- Opne Azure Key Vault i portalen, og gå til bladet Access Policies.
- Klikk på +Add Access Policy.

- Angi nødvendig løyvenivå. I dette tilfellet treng Azure Databricks berre å kunne Get og List secrets som er lagra i Key Vault – ikkje meir og ikkje mindre.

- Til slutt må vi angi principal, altså identiteten til parten som får desse løyva. Då dette blei skrive, fekk ikkje kvart enkelt Databricks-workspace sin eigen 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 sida av Select principal for å velje han. Søk etter AzureDatabricks; det skal berre vere éi oppføring. Vel elementet i lista, klikk på Select, og klikk deretter på Add.

- Du kjem då tilbake til bladet Access Policies. Klikk på Save for å bruke endringane.
Opprette secret scope i Azure Databricks
No som vi har oppretta løyve på Key Vault for den globale AzureDatabricks enterprise application, må vi konfigurere workspacet til å kople til Key Vault og hente secrets ved å opprette eit secret scope.
NB: Du treng minst Contributor-løyve på sjølve Key Vault-en for å utføre denne operasjonen.
- Opne ei ny fane, og gå til Databricks-workspacet ditt. URL-en skal likne på https://adb-..azuredatabricks.net/.
- Legg til #secrets/createScope i URL-en, slik at han blir https://adb-..azuredatabricks.net/#secrets/createScope.
- Merk at denne ekstra delen av URL-en skil mellom store og små bokstavar.
- Gå til Key Vault-en din i Azure Portal i ei anna fane, og opne bladet Properties.
- Noter følgjande verdiar:
- Gå tilbake til Databricks-workspace-fana med sida for å opprette secret scope.
- Angi scope-namnet som verdien du definerte på linje 5 i del 4 av opprettinga av notebooken. I eksempelet bruker eg databricks-secret-scope. I meir realistiske scenario med ulike miljø, som Development, Test, Production og Quality Assurance, vil du truleg ha Key Vault-ar avgrensa til kvart miljø med eit tilsvarande secret scope.
- I eksempelmiljøet mitt krevst ikkje høgaste tryggleiksnivå, så eg lèt alle brukarar administrere principal-en i rullegardinlista nedanfor.
- Lim inn Vault URI i feltet DNS Name.
- Lim inn Resource ID i feltet Resource ID.
- Klikk på Create.
- Viss brukaren som utfører operasjonen har minst Contributor-løyve, skal følgjande melding visast.

Kople Azure Databricks til Azure Data Factory
No 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.
Gje Azure Data Factory løyve på Azure Databricks-workspaceet
- Gå til Azure Databricks-workspacet i Azure Portal i ei ny fane.
- Opne bladet Access control (IAM).
- Klikk på Add.
- Klikk på Add role assignment.

- Klikk på rada med namnet Contributor. Azure Databricks har dessverre ikkje eit eige sett med spesifikke roller for Azure Databricks, så den generiske Contributor-rolla må brukast.
- Klikk på Next.
- På fana Members set du alternativknappen ved sida av Assign access to til Managed identity.
- Klikk på Select members.
- Finn Azure Data Factory-instansen din ved å velje rett subscription og deretter filtrere på managed identity-type. I dette tilfellet er det Data factory (V2).

- Vel rett data factory, og klikk deretter på Select.
- Klikk på Review + assign og deretter Assign for å fullføre prosessen.
Leggje til Azure Databricks som ein Linked Service i Azure Data Factory
- Opne Azure Data Factory studio som du har brukt i denne øvinga, i ei ny fane.
- Klikk på ikonet
. - Klikk på bladet Linked services.
- Klikk på New.
- Klikk på fana Compute, vel Azure Databricks, og klikk deretter på Continue.

- Gi linked service namnet LS_AzureDatabricks, eller eit anna namn etter namnekonvensjonen din.
- Konfigurer han slik:
- Bruk standard Integration runtime (AutoResolveIntegrationRuntime), med mindre du har eit bestemt behov for ein annan.
- Vel From Azure subscription i rullegardinlista for account selection method.
- Vel rett Azure subscription, og deretter rett workspace.
- Vel ein ny job cluster som cluster type for å redusere kostnadene.
- Angi autentiseringstype til Managed service identity.
- Workspace resource ID skal vere førehandsutfylt.
- Vel cluster version 9.1 LTS. Det blei innført ei brytande endring i Spark 3.2 som gjer at koden i steg 6 ikkje fungerer. Bruk Databricks 9.1 LTS Runtime.
- Standard_DS3_v2 skal passe som cluster node type.
- Vel Python version 3.
- Angi autoskalering av workers til minimum 1 og maksimum 2.

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

- Du skal no sjå LS_AzureDatabricks i lista over linked services.

Leggje Azure Databricks-notebooken til i den eksisterande Azure Data Factory-pipelinen
- Gå til author mode i Azure Data Factory studio ved å klikke på ikonet
. - Finn den eksisterande pipelinen som blei oppretta tidlegare, og opne han.
- Utvid alle aktivitetane.

- Dra ein Databricks-aktivitet av typen Notebook inn i pipelinen, og legg han til som etterfølgjande aktivitet etter Copy Raw Data. Kall aktiviteten Transform Source Data.

- Klikk på fana Azure Databricks for å knyte aktiviteten til den nyleg oppretta linked service-en.

- Klikk på fana Settings, og klikk deretter på Browse.
- Søk etter notebooken du lagra tidlegare, og klikk på OK.

- Til slutt må vi angi base parameters for Databricks-aktiviteten, slik at folderPath og pipelineRunId kan sendast til notebooken ved køyring.
- Utvid Base parameters.
- Klikk på +New to gonger.
- Gi den første parameteren namnet filename.
- Angi ein dynamic content-verdi med følgjande kode:
- Gi den andre parameteren namnet pipelineRunId.
- Angi ein dynamic content-verdi med følgjande kode:
- Klikk på knappen Publish all.
- Køyr Debug for pipelinen for å kontrollere at alt fungerer. Ver merksam på at viss kjeldedatasettet er stort, kan Databricks-notebooken bruke nokre minutt på behandling etter at clusteret er oppretta. Eg opplevde ein flaskehals i steget der dataa blir lagra i Delta Lake. Det kan vere verdt å auke talet på workers som er tilgjengelege for job clusteret, men då må Azure subscription ha rett core-kvote for VM Series som blir brukte av Databricks-clusteret. Dei fleste VM Series har som standard 10 cores. Du kan be om kvoteauke ved å sende inn ein support request.
Oppsummering
Då er vi i mål: Vi har laga Databricks-notebooken som behandlar Log Analytics-dataa, og integrert notebooken i den eksisterande Azure Data Factory-pipelinen.