Complete Azure Databricks for ML Tutorial: Unified Analytics Platform
Azure Databricks provides a collaborative Apache Spark-based analytics platform optimized for machine learning. It combines data engineering, data science, and machine learning on a unified platform.
Why Azure Databricks for ML?
Key Benefits:- Unified platform: Data engineering and ML in one place
- Collaborative: Notebooks with real-time collaboration
- Scalable: Auto-scaling Spark clusters
- MLflow integration: Built-in experiment tracking
- Delta Lake: Reliable data lake storage
- Databricks Workspace
- Spark Clusters
- Notebooks
- MLflow
- Feature Store
- Model Serving
Prerequisites
pip install databricks-sdk mlflow
Azure CLI
az login
Setup
1. Create Databricks Workspace
from azure.mgmt.databricks import AzureDatabricksManagementClient
from azure.identity import DefaultAzureCredential
credential = DefaultAzureCredential()
client = AzureDatabricksManagementClient(
credential=credential,
subscriptionid="your-subscription-id"
)
Create workspace
workspace = client.workspaces.begincreateorupdate(
resourcegroupname="my-resource-group",
workspacename="my-databricks-workspace",
parameters={
"location": "eastus",
"sku": {"name": "premium"}
}
).result()
print(f"Workspace created: {workspace.name}")
2. Connect with Databricks SDK
from databricks.sdk import WorkspaceClient
Initialize client
w = WorkspaceClient(
host="https://adb-xxxxx.azuredatabricks.net",
token="dapi-xxxxx"
)
List clusters
clusters = w.clusters.list()
for cluster in clusters:
print(f"{cluster.clustername}: {cluster.state}")
Cluster Management
1. Create ML Cluster
from databricks.sdk.service.compute import (
ClusterSpec,
AutoScale,
AzureAttributes
)
Create cluster
cluster = w.clusters.create(
clustername="ml-cluster",
sparkversion="13.3.x-ml-scala2.12",
nodetypeid="StandardDS3v2",
autoscale=AutoScale(minworkers=1, maxworkers=8),
azureattributes=AzureAttributes(
availability="ONDEMANDAZURE",
firstondemand=1
),
sparkconf={
"spark.databricks.delta.preview.enabled": "true"
},
customtags={
"project": "ml-training",
"team": "data-science"
}
).result()
print(f"Cluster ID: {cluster.clusterid}")
2. GPU Cluster for Deep Learning
gpucluster = w.clusters.create(
cluster
name="gpu-ml-cluster",
sparkversion="13.3.x-gpu-ml-scala2.12",
nodetypeid="StandardNC6sv3",
numworkers=2,
sparkconf={
"spark.task.resource.gpu.amount": "1"
}
).result()
Notebooks and Data
1. Create Notebook
# Create notebook
notebook = w.workspace.mkdirs("/Users/user@company.com/ml-projects")
Import notebook
w.workspace.import(
path="/Users/user@company.com/ml-projects/training",
format="SOURCE",
language="PYTHON",
content=base64.b64encode(notebookcontent.encode()).decode()
)
2. Working with Delta Lake
# In Databricks notebook
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
Read data
df = spark.read.format("csv").option("header", "true").load("dbfs:/data/train.csv")
Write to Delta Lake
df.write.format("delta").mode("overwrite").save("/delta/trainingdata")
Read from Delta Lake
deltadf = spark.read.format("delta").load("/delta/trainingdata")
Create Delta table
spark.sql("""
CREATE TABLE IF NOT EXISTS mldata.trainingdata
USING DELTA
LOCATION '/delta/trainingdata'
""")
3. Feature Engineering with Spark
from pyspark.sql import functions as F
from pyspark.ml.feature import VectorAssembler, StandardScaler
Feature engineering
df = df.withColumn("featureratio", F.col("featurea") / F.col("featureb"))
df = df.withColumn("logfeature", F.log(F.col("featurec") + 1))
Create feature vector
featurecols = ["featurea", "featureb", "featureratio", "logfeature"]
assembler = VectorAssembler(inputCols=featurecols, outputCol="features")
df = assembler.transform(df)
Scale features
scaler = StandardScaler(inputCol="features", outputCol="scaledfeatures")
scalermodel = scaler.fit(df)
df = scalermodel.transform(df)
MLflow Integration
1. Experiment Tracking
import mlflow
from mlflow.tracking import MlflowClient
Set experiment
mlflow.setexperiment("/Users/user@company.com/my-experiment")
Start run
with mlflow.startrun(runname="sklearn-training"):
# Log parameters
mlflow.logparam("modeltype", "RandomForest")
mlflow.logparam("nestimators", 100)
# Train model
from sklearn.ensemble import RandomForestClassifier
model = RandomForestClassifier(nestimators=100)
model.fit(Xtrain, ytrain)
# Log metrics
accuracy = model.score(Xtest, ytest)
mlflow.logmetric("accuracy", accuracy)
# Log model
mlflow.sklearn.logmodel(model, "model")
print(f"Accuracy: {accuracy}")
2. Autologging
# Enable autologging
mlflow.autolog()
Train with Spark ML
from pyspark.ml.classification import RandomForestClassifier
from pyspark.ml.evaluation import MulticlassClassificationEvaluator
rf = RandomForestClassifier(
labelCol="label",
featuresCol="features",
numTrees=100
)
Training automatically logged
model = rf.fit(traindf)
predictions = model.transform(testdf)
evaluator = MulticlassClassificationEvaluator(
labelCol="label",
predictionCol="prediction",
metricName="accuracy"
)
accuracy = evaluator.evaluate(predictions)
print(f"Accuracy: {accuracy}")
3. Model Registry
# Register model
modeluri = f"runs:/{mlflow.activerun().info.runid}/model"
mlflow.registermodel(modeluri, "production-classifier")
Transition model stage
client = MlflowClient()
client.transitionmodelversionstage(
name="production-classifier",
version=1,
stage="Production"
)
Distributed Training
1. Spark ML Pipeline
from pyspark.ml import Pipeline
from pyspark.ml.feature import StringIndexer, VectorAssembler
from pyspark.ml.classification import GBTClassifier
Create pipeline stages
labelindexer = StringIndexer(inputCol="target", outputCol="label")
assembler = VectorAssembler(inputCols=featurecols, outputCol="features")
gbt = GBTClassifier(labelCol="label", featuresCol="features", maxIter=50)
Create pipeline
pipeline = Pipeline(stages=[labelindexer, assembler, gbt])
Train
with mlflow.startrun():
model = pipeline.fit(traindf)
mlflow.spark.logmodel(model, "spark-model")
Predict
predictions = model.transform(testdf)
2. Hyperparameter Tuning
from pyspark.ml.tuning import ParamGridBuilder, CrossValidator
Create parameter grid
paramgrid = ParamGridBuilder() \
.addGrid(gbt.maxDepth, [5, 10, 15]) \
.addGrid(gbt.maxIter, [20, 50, 100]) \
.build()
Create cross validator
cv = CrossValidator(
estimator=pipeline,
estimatorParamMaps=paramgrid,
evaluator=evaluator,
numFolds=3,
parallelism=4
)
Run cross-validation
with mlflow.startrun():
cvmodel = cv.fit(traindf)
# Log best parameters
bestmodel = cvmodel.bestModel
mlflow.logparam("bestmaxdepth", bestmodel.stages[-1].getMaxDepth())
# Log metrics
predictions = cvmodel.transform(testdf)
accuracy = evaluator.evaluate(predictions)
mlflow.logmetric("cvaccuracy", accuracy)
3. Distributed Deep Learning with Horovod
from sparkdl import HorovodRunner
def trainhvd():
import horovod.torch as hvd
import torch
import torch.nn as nn
import torch.optim as optim
hvd.init()
# Set device
device = torch.device("cuda" if torch.cuda.isavailable() else "cpu")
# Create model
model = nn.Sequential(
nn.Linear(10, 50),
nn.ReLU(),
nn.Linear(50, 2)
).to(device)
# Horovod: wrap optimizer
optimizer = optim.Adam(model.parameters(), lr=0.001 hvd.size())
optimizer = hvd.DistributedOptimizer(optimizer)
# Horovod: broadcast parameters
hvd.broadcastparameters(model.statedict(), rootrank=0)
# Training loop
for epoch in range(10):
# ... training code
pass
return model
Run distributed training
hr = HorovodRunner(np=4)
model = hr.run(trainhvd)
Feature Store
1. Create Feature Table
from databricks.featurestore import FeatureStoreClient
fs = FeatureStoreClient()
Create feature table
fs.createtable(
name="mlfeatures.customerfeatures",
primarykeys=["customerid"],
df=featuredf,
description="Customer features for churn prediction"
)
2. Update Features
# Write features
fs.writetable(
name="mlfeatures.customerfeatures",
df=newfeaturesdf,
mode="merge"
)
3. Train with Feature Store
from databricks.featurestore import FeatureLookup
Define feature lookups
featurelookups = [
FeatureLookup(
tablename="mlfeatures.customerfeatures",
featurenames=["tenure", "monthlycharges", "totalcharges"],
lookupkey="customerid"
)
]
Create training set
trainingset = fs.createtrainingset(
df=labeldf,
featurelookups=featurelookups,
label="churn"
)
Train model
trainingdf = trainingset.loaddf()
with mlflow.startrun():
model = trainmodel(trainingdf)
# Log model with feature store
fs.logmodel(
model,
artifactpath="model",
flavor=mlflow.sklearn,
trainingset=trainingset,
registeredmodelname="churn-model"
)
Model Serving
1. Enable Model Serving
from databricks.sdk.service.serving import (
EndpointCoreConfigInput,
ServedModelInput,
ServedModelInputWorkloadSize
)
Create serving endpoint
endpoint = w.servingendpoints.create(
name="churn-prediction-endpoint",
config=EndpointCoreConfigInput(
servedmodels=[
ServedModelInput(
modelname="churn-model",
modelversion="1",
workloadsize=ServedModelInputWorkloadSize.SMALL,
scaletozeroenabled=True
)
]
)
)
print(f"Endpoint URL: {endpoint.url}")
2. Query Endpoint
import requests
Get endpoint URL
endpoint = w.servingendpoints.get("churn-prediction-endpoint")
Query endpoint
response = requests.post(
f"{endpoint.config.servedmodels[0].endpointurl}/invocations",
headers={
"Authorization": f"Bearer {token}",
"Content-Type": "application/json"
},
json={
"dataframerecords": [
{"customerid": "123", "tenure": 12, "monthlycharges": 50.0}
]
}
)
print(response.json())
Jobs and Workflows
1. Create ML Job
from databricks.sdk.service.jobs import (
JobSettings,
Task,
NotebookTask,
JobCluster
)
Create job
job = w.jobs.create(
name="daily-model-training",
tasks=[
Task(
taskkey="trainmodel",
notebooktask=NotebookTask(
notebookpath="/Users/user@company.com/ml-projects/training"
),
jobclusterkey="ml-cluster"
)
],
jobclusters=[
JobCluster(
jobclusterkey="ml-cluster",
newcluster={
"sparkversion": "13.3.x-ml-scala2.12",
"nodetypeid": "StandardDS3v2",
"numworkers": 2
}
)
],
schedule={
"quartzcronexpression": "0 0 6 * ?",
"timezoneid": "UTC"
}
)
print(f"Job ID: {job.jobid}")
2. Run Job
# Run job
run = w.jobs.runnow(jobid=job.jobid)
print(f"Run ID: {run.runid}")
Wait for completion
runstate = w.jobs.getrun(runid=run.runid)
print(f"State: {runstate.state.lifecyclestate}")
Best Practices
1. Cluster Configuration
# Use autoscaling for variable workloads
clusterconfig = {
"autoscale": {"minworkers": 1, "maxworkers": 10},
"sparkversion": "13.3.x-ml-scala2.12",
"nodetypeid": "StandardDS3v2",
"sparkconf": {
"spark.databricks.io.cache.enabled": "true",
"spark.sql.shuffle.partitions": "auto"
}
}
2. Delta Lake Optimization
# Optimize Delta table
spark.sql("OPTIMIZE mldata.trainingdata ZORDER BY (customerid)")
Vacuum old files
spark.sql("VACUUM mldata.trainingdata RETAIN 168 HOURS")
Enable auto-optimization
spark.conf.set("spark.databricks.delta.optimizeWrite.enabled", "true")
spark.conf.set("spark.databricks.delta.autoCompact.enabled", "true")
Conclusion
Azure Databricks for ML provides:
Key takeaways:
- Use Spark for large-scale data processing
- Leverage MLflow for experiment tracking
- Use Feature Store for feature management
- Deploy models with Model Serving
- Optimize Delta Lake for performance