Espandere le funzionalità di Spark con MLflow

Ciao, membri di Habr. Come già scritto, questo mese OTUS lancia due corsi di machine learning, precisamente base e avanzato. Pertanto, continuiamo a condividere materiale utile.

L'obiettivo di questo articolo è raccontare della nostra prima esperienza nell'utilizzo di MLflow.

. Inizieremo a rivedere MLflow il suo server di tracciamento e loggheremo tutte le iterazioni dello studio. Poi condivideremo l'esperienza di collegare Spark con MLflow tramite UDF.

Contesto

Noi di Alpha Health utilizza il machine learning e l'intelligenza artificiale per permettere alle persone di prendersi cura della propria salute e benessere. Per questo motivo, i modelli di machine learning sono alla base dei prodotti di trattamento dei dati che sviluppiamo, ed è per questo che siamo stati attratti da MLflow — una piattaforma open source che copre tutti gli aspetti del ciclo di vita del machine learning.

MLflow

L'obiettivo principale di MLflow è fornire uno strato addizionale sopra il machine learning, che consentirebbe agli specialisti del data science di lavorare praticamente con qualsiasi libreria di machine learning (h2o, keras, mleap, pytorch, sklearn e tensorflow), portando il loro lavoro a un nuovo livello.

MLflow offre tre componenti:

  • Tracking – registrazione e interrogazione degli esperimenti: codice, dati, configurazione e risultati. Monitorare il processo di creazione del modello è fondamentale.
  • Progetti – Formato per l'imballaggio da eseguire su qualsiasi piattaforma (ad esempio, SageMaker)
  • Models – formato comune per inviare modelli a vari strumenti di distribuzione.

MLflow (al momento della scrittura dell'articolo in versione alpha) è una piattaforma open source che consente di gestire il ciclo di vita del machine learning, inclusi esperimenti, riutilizzo e distribuzione.

Configurazione di MLflow

Per utilizzare MLflow è necessario prima configurare tutto l'ambiente Python, per questo utilizzeremo PyEnv (per installare Python su Mac, consulta qui). Così possiamo creare un ambiente virtuale dove installeremo tutte le librerie necessarie per eseguire il software.

```
pyenv install 3.7.0
pyenv global 3.7.0 # Usa Python 3.7
mkvirtualenv mlflow # Crea un Ambiente Virtuale con Python 3.7
workon mlflow
```

Installeremo le librerie richieste.

```
pip install mlflow==0.7.0 
            Cython==0.29  
            numpy==1.14.5 
            pandas==0.23.4 
            pyarrow==0.11.0
```

Nota: utilizziamo PyArrow per eseguire modelli come UDF. Le versioni di PyArrow e Numpy dovevano essere corrette, poiché le ultime versioni presentavano conflitti tra loro.

Avviamo l'interfaccia utente di Tracking

MLflow Tracking ci consente di registrare e interrogare esperimenti utilizzando Python e REST API. Inoltre, è possibile specificare dove archiviare gli artefatti del modello (localhost, Amazon S3, Azure Blob Storage, Google Cloud Storage o server SFTP). Poiché in Alpha Health utilizziamo AWS, S3 sarà utilizzato come archivio per gli artefatti.

# Running a Tracking Server
mlflow server 
    --file-store /tmp/mlflow/fileStore 
    --default-artifact-root s3://<bucket>/mlflow/artifacts/ 
    --host localhost
    --port 5000

MLflow consiglia di utilizzare un'archiviazione di file persistente. L'archiviazione di file è il luogo in cui il server memorizzerà i metadati delle esecuzioni e degli esperimenti. Quando avvii il server, assicurati che punti a un'archiviazione di file persistente. Qui per l'esperimento utilizzeremo semplicemente /tmp.

Ricorda che, se desideriamo utilizzare il server mlflow per eseguire esperimenti precedenti, questi devono essere presenti nell'archiviazione di file. Tuttavia, anche senza questo, potremmo usarli in UDF, poiché abbiamo solo bisogno del percorso del modello.

Nota: Tieni presente che Tracking UI e il client del modello devono avere accesso alla posizione dell'artefatto. Ciò significa che, indipendentemente dal fatto che Tracking UI si trovi in un'istanza EC2, quando si esegue MLflow localmente, la macchina deve avere accesso diretto a S3 per scrivere modelli artefatti.

Espandere le funzionalità di Spark con MLflow
La Tracking UI memorizza gli artefatti in un bucket S3

Esecuzione di modelli

Una volta che il server di Tracking è in funzione, puoi iniziare ad addestrare i modelli.

Come esempio, utilizzeremo la modifica wine dall'esempio di MLflow in Sklearn.

MLFLOW_TRACKING_URI=http://localhost:5000 python wine_quality.py 
  --alpha 0.9
  --l1_ratio 0.5
  --wine_file ./data/winequality-red.csv

Come già accennato, MLflow consente di registrare parametri, metriche e artefatti dei modelli, in modo da poter monitorare come si sviluppano nel corso delle iterazioni. Questa funzione è estremamente utile, in quanto ci consente di riprodurre il miglior modello, facendo riferimento al server di Tracking o comprendendo quale codice ha eseguito l'iterazione desiderata, utilizzando i log degli hash commit di git.

with mlflow.start_run():

    ... modello ...

    mlflow.log_param("source", wine_path)
    mlflow.log_param("alpha", alpha)
    mlflow.log_param("l1_ratio", l1_ratio)

    mlflow.log_metric("rmse", rmse)
    mlflow.log_metric("r2", r2)
    mlflow.log_metric("mae", mae)

    mlflow.set_tag('domain', 'wine')
    mlflow.set_tag('predict', 'quality')
    mlflow.sklearn.log_model(lr, "model")

Espandere le funzionalità di Spark con MLflow
Iterazioni wine

Parte server per il modello

Il server di tracciamento MLflow, avviato con il comando “mlflow server”, ha un'API REST per monitorare le esecuzioni e registrare i dati nel filesystem locale. Puoi specificare l'indirizzo del server di tracciamento utilizzando la variabile d'ambiente «MLFLOW_TRACKING_URI» e l'API di tracciamento MLflow si connetterà automaticamente al server di tracciamento a quell'indirizzo per creare/ottenere informazioni sull'esecuzione, metriche dei log, ecc.

Fonte: Docs// Esecuzione di un server di tracciamento

Per fornire il modello con un server, avremo bisogno di un server di tracciamento avviato (vedi l'interfaccia di avvio) e dell'Run ID del modello.

Espandere le funzionalità di Spark con MLflow
Run ID

# Serve a sklearn model through 127.0.0.0:5005
MLFLOW_TRACKING_URI=http://0.0.0.0:5000 mlflow sklearn serve 
  --port 5005  
  --run_id 0f8691808e914d1087cf097a08730f17 
  --model-path model

Per gestire i modelli tramite la funzionalità MLflow serve, avremo bisogno di accesso all'interfaccia utente di tracciamento per ottenere informazioni sul modello semplicemente specificando --run_id.

Una volta che il modello è connesso al server di tracciamento, possiamo ottenere un nuovo endpoint per il modello.

# Query Tracking Server Endpoint
curl -X POST 
  http://127.0.0.1:5005/invocations 
  -H 'Content-Type: application/json' 
  -d '[
	{
		"fixed acidity": 3.42, 
		"volatile acidity": 1.66, 
		"citric acid": 0.48, 
		"residual sugar": 4.2, 
		"chloridessssss": 0.229, 
		"free sulfur dsioxide": 19, 
		"total sulfur dioxide": 25, 
		"density": 1.98, 
		"pH": 5.33, 
		"sulphates": 4.39, 
		"alcohol": 10.8
	}
]'

> {"predictions": [5.825055635303461]}

Esecuzione di modelli da Spark

Sebbene il server di tracciamento sia abbastanza potente per gestire modelli in tempo reale, il suo utilizzo per l'addestramento e la funzionalità serve (fonte: mlflow // docs // models # local), l'applicazione di Spark (batch o streaming) è una soluzione ancora più potente grazie alla distribuzione.

Immagina di aver appena eseguito l'addestramento in offline e poi di aver applicato il modello risultante a tutti i tuoi dati. È qui che Spark e MLflow si distinguono.

Installiamo PySpark + Jupyter + Spark

Fonte: Iniziare con PySpark — Jupyter

Per mostrare come applichiamo i modelli MLflow ai DataFrame Spark, è necessario configurare la collaborazione tra i Jupyter notebooks e PySpark.

Inizia installando l'ultima versione stabile Apache Spark:

cd ~/Downloads/
tar -xzf spark-2.4.3-bin-hadoop2.7.tgz
mv ~/Downloads/spark-2.4.3-bin-hadoop2.7 ~/
ln -s ~/spark-2.4.3-bin-hadoop2.7 ~/spark̀

Installa PySpark e Jupyter in un ambiente virtuale:

pip install pyspark jupyter

Configura le variabili d'ambiente:

export SPARK_HOME=~/spark
export PATH=$SPARK_HOME/bin:$PATH
export PYSPARK_DRIVER_PYTHON=jupyter
export PYSPARK_DRIVER_PYTHON_OPTS="notebook --notebook-dir=${HOME}/Projects/notebooks"

Definendo notebook-dir, saremo in grado di memorizzare i nostri notebook nella cartella desiderata.

Avviamo Jupyter da PySpark

Poiché siamo riusciti a configurare Jupiter come driver PySpark, ora possiamo avviare il Jupyter notebook nel contesto di PySpark.

(mlflow) afranzi:~$ pyspark
[I 19:05:01.572 NotebookApp] estensione sparkmagic abilitata!
[I 19:05:01.573 NotebookApp] Servendo notebook dalla directory locale: /Users/afranzi/Projects/notebooks
[I 19:05:01.573 NotebookApp] Il Jupyter Notebook è in esecuzione su:
[I 19:05:01.573 NotebookApp] http://localhost:8888/?token=c06252daa6a12cfdd33c1d2e96c8d3b19d90e9f6fc171745
[I 19:05:01.573 NotebookApp] Usa Control-C per fermare questo server e chiudere tutti i kernel (due volte per saltare la conferma).
[C 19:05:01.574 NotebookApp]

    Copia/incolla questo URL nel tuo browser quando ti connetti per la prima volta,
    per accedere con un token:
        http://localhost:8888/?token=c06252daa6a12cfdd33c1d2e96c8d3b19d90e9f6fc171745

Espandere le funzionalità di Spark con MLflow

Come accennato in precedenza, MLflow fornisce una funzione di registrazione degli artefatti del modello in S3. Una volta che abbiamo in mano il modello scelto, abbiamo la possibilità di importarlo come UDF tramite il modulo mlflow.pyfunc.

import mlflow.pyfunc

model_path = 's3:///mlflow/artifacts/1/0f8691808e914d1087cf097a08730f17/artifacts/model'
wine_path = '/Users/afranzi/Projects/data/winequality-red.csv'
wine_udf = mlflow.pyfunc.spark_udf(spark, model_path)

df = spark.read.format("csv").option("header", "true").option('delimiter', ';').load(wine_path)
columns = [ "fixed acidity", "volatile acidity", "citric acid",
            "residual sugar", "chlorides", "free sulfur dioxide",
            "total sulfur dioxide", "density", "pH",
            "sulphates", "alcohol"
          ]
          
df.withColumn('prediction', wine_udf(*columns)).show(100, False)

Espandere le funzionalità di Spark con MLflow
PySpark – Risultato della previsione sulla qualità del vino

Fino a questo momento abbiamo parlato di come utilizzare PySpark con MLflow, eseguendo previsioni sulla qualità del vino su tutto il set di dati wine. Ma cosa fare se dobbiamo utilizzare i moduli Python MLflow da Scala Spark?

Abbiamo testato anche questo, condividendo il contesto Spark tra Scala e Python. Cioè, abbiamo registrato l'UDF MLflow in Python e lo abbiamo utilizzato da Scala (sì, può non essere la migliore soluzione, ma è ciò che abbiamo).

Scala Spark + MLflow

Per questo esempio, aggiungeremo Toree Kernel nel già esistente Jupyter.

Installiamo Spark + Toree + Jupyter

pip install toree
jupyter toree install --spark_home=${SPARK_HOME} --sys-prefix
jupyter kernelspec list
```
```
Kernel disponibili:
  apache_toree_scala    /Users/afranzi/.virtualenvs/mlflow/share/jupyter/kernels/apache_toree_scala
  python3               /Users/afranzi/.virtualenvs/mlflow/share/jupyter/kernels/python3
```

Come si può vedere dal notebook allegato, l'UDF viene utilizzato congiuntamente in Spark e PySpark. Speriamo che questa parte sia utile a coloro che amano Scala e vogliono implementare modelli di machine learning in produzione.

import org.apache.spark.sql.functions.col
import org.apache.spark.sql.types.StructType
import org.apache.spark.sql.{Column, DataFrame}
import scala.util.matching.Regex

val FirstAtRe: Regex = "^_".r
val AliasRe: Regex = "[\s_.:@]+".r

def getFieldAlias(field_name: String): String = {
    FirstAtRe.replaceAllIn(AliasRe.replaceAllIn(field_name, "_"), "")
}

def selectFieldsNormalized(columns: List[String])(df: DataFrame): DataFrame = {
    val fieldsToSelect: List[Column] = columns.map(field =>
        col(field).as(getFieldAlias(field))
    )
    df.select(fieldsToSelect: _*)
}

def normalizeSchema(df: DataFrame): DataFrame = {
    val schema = df.columns.toList
    df.transform(selectFieldsNormalized(schema))
}

FirstAtRe = ^_
AliasRe = [s_.:@]+

getFieldAlias: (field_name: String)String
selectFieldsNormalized: (columns: List[String])(df: org.apache.spark.sql.DataFrame)org.apache.spark.sql.DataFrame
normalizeSchema: (df: org.apache.spark.sql.DataFrame)org.apache.spark.sql.DataFrame
Out[1]:
[s_.:@]+
In [2]:
val winePath = "~\/Research\/mlflow-workshop\/examples\/wine_quality\/data\/winequality-red.csv"
val modelPath = "\/tmp\/mlflow\/artifactStore\/0\/96cba14c6e4b452e937eb5072467bf79\/artifacts\/model"

winePath = ~\/Research\/mlflow-workshop\/examples\/wine_quality\/data\/winequality-red.csv
modelPath = \/tmp\/mlflow\/artifactStore\/0\/96cba14c6e4b452e937eb5072467bf79\/artifacts\/model
Out[2]:
\/tmp\/mlflow\/artifactStore\/0\/96cba14c6e4b452e937eb5072467bf79\/artifacts\/model
In [3]:
val df = spark.read
              .format("csv")
              .option("header", "true")
              .option("delimiter", ";")
              .load(winePath)
              .transform(normalizeSchema)

df = [fixed_acidity: string, volatile_acidity: string ... 10 more fields]
Out[3]:
[fixed_acidity: string, volatile_acidity: string ... 10 more fields]
In [4]:
%%PySpark
import mlflow
from mlflow import pyfunc

model_path = "\/tmp\/mlflow\/artifactStore\/0\/96cba14c6e4b452e937eb5072467bf79\/artifacts\/model"
wine_quality_udf = mlflow.pyfunc.spark_udf(spark, model_path)

spark.udf.register("wineQuality", wine_quality_udf)
Out[4]:
<function spark_udf..predict at 0x1116a98c8>
In [6]:
df.createOrReplaceTempView("wines")
In [10]:
%%SQL
SELECT 
    quality,
    wineQuality(
        fixed_acidity,
        volatile_acidity,
        citric_acid,
        residual_sugar,
        chlorides,
        free_sulfur_dioxide,
        total_sulfur_dioxide,
        density,
        pH,
        sulphates,
        alcohol
    ) AS prediction
FROM wines
LIMIT 10
Out[10]:
+-------+------------------+
|quality|        prediction|
+-------+------------------+
|      5| 5.576883967129615|
|      5|  5.50664776916154|
|      5| 5.525504822954496|
|      6| 5.504311247097457|
|      5| 5.576883967129615|
|      5|5.5556903912725755|
|      5| 5.467882654744997|
|      7| 5.710602976324739|
|      7| 5.657319539336507|
|      5| 5.345098606538708|
+-------+------------------+

In [17]:
spark.catalog.listFunctions.filter('name like "%wineQuality%").show(20, false)

+-----------+--------+-----------+---------+-----------+
|name       |database|description|className|isTemporary|
+-----------+--------+-----------+---------+-----------+
|wineQuality|null    |null       |null     |true       |
+-----------+--------+-----------+---------+-----------+

Passaggi successivi

Nonostante al momento della scrittura di questo articolo MLflow si trovi in versione Alpha, appare già piuttosto promettente. La possibilità di eseguire più framework di machine learning e utilizzarli da un unico punto finale porta i sistemi di raccomandazione a un nuovo livello.

Inoltre, MLflow avvicina gli ingegneri dei dati e gli specialisti di Data Science, creando uno strato comune tra di loro.

Dopo questa ricerca su MLflow, siamo sicuri che andremo avanti e lo utilizzeremo per i nostri pipeline Spark e nei sistemi di raccomandazione.

Sarebbe utile sincronizzare lo storage dei file con il database, anziché con il sistema di file. In questo modo dovremmo ottenere diversi endpoint che possono utilizzare lo stesso storage di file. Ad esempio, utilizzare più istanze Presto e Athena con lo stesso Glue metastore.

In sintesi, vogliamo ringraziare la comunità di MLFlow per rendere il nostro lavoro con i dati più interessante.

Se stai sperimentando con MLflow, non esitare a scriverci e raccontarci come lo stai utilizzando, e ancor di più se lo stai usando in produzione.

Scopri di più sui corsi:
Machine Learning. Corso base
Machine Learning. Corso avanzato

Leggi anche:

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server 🔥 Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server | ProHoster