Tutorial Lengkap Azure Databricks untuk ML: Platform Analytics Terpadu
Azure Databricks menyediakan platform analytics kolaboratif berbasis Apache Spark yang dioptimasi untuk machine learning. Platform ini menggabungkan data engineering, data science, dan machine learning dalam satu platform terpadu.
Mengapa Azure Databricks untuk ML?
Manfaat Utama:- Platform terpadu: Data engineering dan ML dalam satu tempat
- Kolaboratif: Notebooks dengan kolaborasi real-time
- Scalable: Auto-scaling Spark clusters
- Integrasi MLflow: Built-in experiment tracking
- Delta Lake: Storage data lake yang reliable
- Databricks Workspace
- Spark Clusters
- Notebooks
- MLflow
- Feature Store
- Model Serving
Prerequisites
pip install databricks-sdk mlflow
Azure CLI
az login
Setup
1. Buat Databricks Workspace
from azure.mgmt.databricks import AzureDatabricksManagementClient
from azure.identity import DefaultAzureCredential
credential = DefaultAzureCredential()
client = AzureDatabricksManagementClient(
credential=credential,
subscriptionid="your-subscription-id"
)
Buat workspace
workspace = client.workspaces.begincreateorupdate(
resourcegroupname="my-resource-group",
workspacename="my-databricks-workspace",
parameters={
"location": "eastus",
"sku": {"name": "premium"}
}
).result()
print(f"Workspace dibuat: {workspace.name}")
2. Koneksi dengan 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}")
Manajemen Cluster
1. Buat ML Cluster
from databricks.sdk.service.compute import (
ClusterSpec,
AutoScale,
AzureAttributes
)
Buat 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 untuk 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 dan Data
1. Buat Notebook
# Buat 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. Bekerja dengan Delta Lake
# Di Databricks notebook
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
Baca data
df = spark.read.format("csv").option("header", "true").load("dbfs:/data/train.csv")
Tulis ke Delta Lake
df.write.format("delta").mode("overwrite").save("/delta/trainingdata")
Baca dari Delta Lake
deltadf = spark.read.format("delta").load("/delta/trainingdata")
Buat Delta table
spark.sql("""
CREATE TABLE IF NOT EXISTS mldata.trainingdata
USING DELTA
LOCATION '/delta/trainingdata'
""")
3. Feature Engineering dengan 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))
Buat 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)
Integrasi MLflow
1. Experiment Tracking
import mlflow
from mlflow.tracking import MlflowClient
Set experiment
mlflow.setexperiment("/Users/user@company.com/my-experiment")
Mulai 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
# Aktifkan autologging
mlflow.autolog()
Train dengan Spark ML
from pyspark.ml.classification import RandomForestClassifier
from pyspark.ml.evaluation import MulticlassClassificationEvaluator
rf = RandomForestClassifier(
labelCol="label",
featuresCol="features",
numTrees=100
)
Training dilog otomatis
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")
Transisi 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
Buat pipeline stages
labelindexer = StringIndexer(inputCol="target", outputCol="label")
assembler = VectorAssembler(inputCols=featurecols, outputCol="features")
gbt = GBTClassifier(labelCol="label", featuresCol="features", maxIter=50)
Buat 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
Buat parameter grid
paramgrid = ParamGridBuilder() \
.addGrid(gbt.maxDepth, [5, 10, 15]) \
.addGrid(gbt.maxIter, [20, 50, 100]) \
.build()
Buat cross validator
cv = CrossValidator(
estimator=pipeline,
estimatorParamMaps=paramgrid,
evaluator=evaluator,
numFolds=3,
parallelism=4
)
Jalankan 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 dengan 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")
# Buat 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):
# ... kode training
pass
return model
Jalankan distributed training
hr = HorovodRunner(np=4)
model = hr.run(trainhvd)
Feature Store
1. Buat Feature Table
from databricks.featurestore import FeatureStoreClient
fs = FeatureStoreClient()
Buat feature table
fs.createtable(
name="mlfeatures.customerfeatures",
primarykeys=["customerid"],
df=featuredf,
description="Customer features untuk prediksi churn"
)
2. Update Features
# Tulis features
fs.writetable(
name="mlfeatures.customerfeatures",
df=newfeaturesdf,
mode="merge"
)
3. Train dengan Feature Store
from databricks.featurestore import FeatureLookup
Definisikan feature lookups
featurelookups = [
FeatureLookup(
tablename="mlfeatures.customerfeatures",
featurenames=["tenure", "monthlycharges", "totalcharges"],
lookupkey="customerid"
)
]
Buat training set
trainingset = fs.createtrainingset(
df=labeldf,
featurelookups=featurelookups,
label="churn"
)
Train model
trainingdf = trainingset.loaddf()
with mlflow.startrun():
model = trainmodel(trainingdf)
# Log model dengan feature store
fs.logmodel(
model,
artifactpath="model",
flavor=mlflow.sklearn,
trainingset=trainingset,
registeredmodelname="churn-model"
)
Model Serving
1. Aktifkan Model Serving
from databricks.sdk.service.serving import (
EndpointCoreConfigInput,
ServedModelInput,
ServedModelInputWorkloadSize
)
Buat 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
Dapatkan 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 dan Workflows
1. Buat ML Job
from databricks.sdk.service.jobs import (
JobSettings,
Task,
NotebookTask,
JobCluster
)
Buat 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": "Asia/Jakarta"
}
)
print(f"Job ID: {job.jobid}")
2. Jalankan Job
# Jalankan job
run = w.jobs.runnow(jobid=job.jobid)
print(f"Run ID: {run.runid}")
Tunggu selesai
runstate = w.jobs.getrun(runid=run.runid)
print(f"State: {runstate.state.lifecyclestate}")
Best Practices
1. Konfigurasi Cluster
# Gunakan autoscaling untuk workloads variabel
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. Optimasi Delta Lake
# Optimize Delta table
spark.sql("OPTIMIZE mldata.trainingdata ZORDER BY (customerid)")
Vacuum old files
spark.sql("VACUUM mldata.trainingdata RETAIN 168 HOURS")
Aktifkan auto-optimization
spark.conf.set("spark.databricks.delta.optimizeWrite.enabled", "true")
spark.conf.set("spark.databricks.delta.autoCompact.enabled", "true")
Kesimpulan
Azure Databricks untuk ML menyediakan:
Key takeaways:
- Gunakan Spark untuk processing data skala besar
- Manfaatkan MLflow untuk experiment tracking
- Gunakan Feature Store untuk manajemen fitur
- Deploy models dengan Model Serving
- Optimize Delta Lake untuk performa