Erweiterung der Spark-Funktionen mit MLflow

Hallo, Habr-Community. Wie bereits erwähnt, startet OTUS in diesem Monat gleich zwei Kurse zum Thema maschinelles Lernen, nämlich Basislevel und Fortgeschrittenenlevel. In diesem Zusammenhang teilen wir weiterhin nützliche Materialien.

Ziel dieses Artikels ist es, von unserer ersten Erfahrung mit MLflow.

zu berichten. MLflow Wir werden mit dem Tracking-Server beginnen und alle Forschungsiterationen protokollieren. Anschließend teilen wir unsere Erfahrungen mit der Verbindung von Spark und MLflow über UDF.

Kontext

Bei Alpha Health nutzen wir maschinelles Lernen und künstliche Intelligenz, um Menschen zu befähigen, für ihre Gesundheit und ihr Wohlbefinden zu sorgen. Daher sind maschinelle Lernmodelle der Kern unserer entwickelten Datenverarbeitungsprodukte, und genau aus diesem Grund hat MLflow unser Interesse geweckt – eine Open-Source-Plattform, die alle Aspekte des Lebenszyklus des maschinellen Lernens abdeckt.

MLflow

Das Hauptziel von MLflow ist es, eine zusätzliche Schicht über das maschinelle Lernen zu schaffen, die es Data-Science-Experten ermöglicht, praktisch mit jeder Bibliothek für maschinelles Lernen zu arbeiten (h2o, keras, mleap, pytorch, sklearn und tensorflow), um deren Leistung auf ein neues Niveau zu heben.

MLflow bietet drei Komponenten:

  • Tracking – Protokollierung und Abfragen von Experimenten: Code, Daten, Konfiguration und Ergebnisse. Den Modellierungsprozess genau zu verfolgen, ist von großer Bedeutung.
  • Projects – Ein Verpackungsformat für die Ausführung auf jeder Plattform (zum Beispiel, SageMaker)
  • Models – Ein gemeinsames Format zur Übermittlung von Modellen an verschiedene Bereitstellungstools.

MLflow (zum Zeitpunkt des Schreibens in der Alpha-Phase) ist eine Open-Source-Plattform, die es ermöglicht, den Lebenszyklus des maschinellen Lernens zu verwalten, einschließlich Experimente, Wiederverwendbarkeit und Bereitstellung.

Einrichtung von MLflow

Um MLflow zu verwenden, müssen Sie zunächst die gesamte Python-Umgebung einrichten. Dazu verwenden wir PyEnv (um Python auf dem Mac zu installieren, schauen Sie bitte hierhin). So können wir eine virtuelle Umgebung erstellen, in der wir alle notwendigen Bibliotheken installieren.

```
pyenv install 3.7.0
pyenv global 3.7.0 # Verwenden Sie Python 3.7
mkvirtualenv mlflow # Erstellen Sie eine virtuelle Umgebung mit Python 3.7
workon mlflow
```

Installieren wir die erforderlichen Bibliotheken.

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

Hinweis: Wir verwenden PyArrow, um Modelle wie UDF auszuführen. Die Versionen von PyArrow und Numpy mussten angepasst werden, da die neuesten Versionen miteinander in Konflikt gerieten.

Tracking UI starten

MLflow Tracking ermöglicht es uns, Experimente mit Python zu protokollieren und abzufragen sowie REST API. Darüber hinaus kann festgelegt werden, wo die Artefakte des Modells gespeichert werden (localhost, Amazon S3, Azure Blob Storage, Google Cloud Storage oder SFTP-Server). Da wir bei Alpha Health AWS nutzen, wird S3 als Speicherort für die Artefakte verwendet.

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

MLflow empfiehlt die Verwendung eines permanenten Speichers. Der Speicher ist der Ort, an dem der Server Metadaten zu Ausführungen und Experimenten aufbewahrt. Stellen Sie beim Start des Servers sicher, dass er auf einen permanenten Speicher verweist. Hier verwenden wir für das Experiment einfach /tmp.

Bitte beachten Sie, dass, wenn wir den MLflow-Server verwenden möchten, um alte Experimente auszuführen, diese im Speicher vorhanden sein müssen. Dennoch könnten wir sie auch ohne dieses Vorhandensein in UDF verwenden, da wir nur den Pfad zum Modell benötigen.

Hinweis: Bitte beachten Sie, dass die Tracking UI und der Client des Modells Zugriff auf den Standort des Artefakts haben müssen. Das bedeutet, dass unabhängig davon, ob sich die Tracking UI in einer EC2-Instanz befindet, die Maschine beim lokalen Start von MLflow direkten Zugriff auf S3 benötigt, um Modelle als Artefakte schreiben zu können.

Erweiterung der Spark-Funktionen mit MLflow
Die Tracking UI speichert Artefakte in einem S3-Bucket.

Modelle ausführen

Sobald der Tracking-Server funktioniert, können Sie mit dem Training der Modelle beginnen.

Als Beispiel verwenden wir die Modifikation aus dem Wein-Beispiel von MLflow in Sklearn.

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

Wie bereits erwähnt, ermöglicht MLflow das Protokollieren von Parametern, Metriken und Artefakten von Modellen, um nachverfolgen zu können, wie sie sich über Iterationen entwickeln. Diese Funktion ist äußerst nützlich, da wir so das beste Modell reproduzieren können, indem wir den Tracking-Server konsultieren oder verstehen, welcher Code die gewünschte Iteration ausgeführt hat, indem wir die Logs der Git-Hash-Commits nutzen.

mit mlflow.start_run():

    ... model ...

    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', 'wein')
    mlflow.set_tag('predict', 'qualität')
    mlflow.sklearn.log_model(lr, "modell")

Erweiterung der Spark-Funktionen mit MLflow
Iterationen Wein

Backend für das Modell

Der MLflow-Tracking-Server, der mit dem Befehl „mlflow server“ gestartet wurde, verfügt über eine REST-API zur Verfolgung von Ausführungen und zur Speicherung von Daten im lokalen Dateisystem. Sie können die Adresse des Tracking-Servers über die Umgebungsvariable „MLFLOW_TRACKING_URI“ angeben, und die Tracking-API von MLflow verbindet sich automatisch mit dem Tracking-Server unter dieser Adresse, um Informationen über Ausführungen, Log-Metriken usw. zu erstellen/abzurufen.

Quelle: Docs// Ausführen eines Tracking-Servers

Um das Modell mit einem Server zu versorgen, benötigen wir einen laufenden Tracking-Server (siehe Start-Schnittstelle) und die Run-ID des Modells.

Erweiterung der Spark-Funktionen mit 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

Für die Bereitstellung von Modellen mit der Funktionalität von MLflow serve benötigen wir Zugriff auf die Tracking-UI, um die Informationen zum Modell einfach anzugeben --run_id.

Sobald das Modell mit dem Tracking-Server verbunden ist, können wir einen neuen Endpunkt für das Modell erhalten.

# 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]}

Ausführen von Modellen aus Spark

Obwohl der Tracking-Server leistungsstark genug ist, um Modelle in Echtzeit zu bedienen, deren Schulung zu ermöglichen und die Funktionalität zu nutzen (Quelle: mlflow // docs // models # local), ist die Verwendung von Spark (Batch oder Streaming) eine noch leistungsfähigere Lösung aufgrund der Verteilung.

Stellen Sie sich vor, Sie hätten ein Training offline durchgeführt, und wenden dann das ausgegebene Modell auf all Ihre Daten an. Genau hier werden Spark und MLflow ihre Stärken zeigen.

Installieren Sie PySpark + Jupyter + Spark

Quelle: Erste Schritte mit PySpark — Jupyter

Um zu zeigen, wie wir MLflow-Modelle auf Spark-Datenframes anwenden, müssen wir die Zusammenarbeit von Jupyter-Notebooks mit PySpark einrichten.

Beginnen Sie mit der Installation der neuesten stabilen Version 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

Installieren Sie PySpark und Jupyter in einer virtuellen Umgebung:

pip install pyspark jupyter

Richten Sie Umgebungsvariablen ein:

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"

Durch die Definition von notebook-dir, können wir unsere Notebooks im gewünschten Ordner speichern.

Starten Sie Jupyter aus PySpark

Nachdem wir Jupiter als PySpark-Treiber konfiguriert haben, können wir jetzt ein Jupyter-Notebook im Kontext von PySpark ausführen.

(mlflow) afranzi:~$ pyspark
[I 19:05:01.572 NotebookApp] sparkmagic-Erweiterung aktiviert!
[I 19:05:01.573 NotebookApp] Notebooks werden aus dem lokalen Verzeichnis bereitgestellt: /Users/afranzi/Projects/notebooks
[I 19:05:01.573 NotebookApp] Das Jupyter-Notebook läuft unter:
[I 19:05:01.573 NotebookApp] http://localhost:8888/?token=c06252daa6a12cfdd33c1d2e96c8d3b19d90e9f6fc171745
[I 19:05:01.573 NotebookApp] Verwenden Sie Control-C, um diesen Server zu stoppen und alle Kerne herunterzufahren (zweimal, um die Bestätigung zu überspringen).
[C 19:05:01.574 NotebookApp]

    Kopieren Sie diesen URL in Ihren Browser, wenn Sie sich zum ersten Mal verbinden,
    um sich mit einem Token einzuloggen:
        http://localhost:8888/?token=c06252daa6a12cfdd33c1d2e96c8d3b19d90e9f6fc171745

Erweiterung der Spark-Funktionen mit MLflow

Wie bereits erwähnt, bietet MLflow eine Funktion zum Protokollieren von Modellartefakten in S3. Sobald wir das gewählte Modell zur Verfügung haben, können wir es mit dem Modul 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)

Erweiterung der Spark-Funktionen mit MLflow
PySpark – Vorhersage der Weinqualität

Bis zu diesem Punkt haben wir darüber gesprochen, wie man PySpark mit MLflow nutzt, um die Weinqualität auf dem gesamten Datensatz 'wine' vorherzusagen. Aber was tun, wenn man MLflow-Module in Scala Spark verwenden möchte?

Wir haben dies getestet, indem wir den Spark-Kontext zwischen Scala und Python aufgeteilt haben. Das heißt, wir haben ein MLflow UDF in Python registriert und es aus Scala verwendet (ja, es mag nicht die beste Lösung sein, aber das haben wir).

Scala Spark + MLflow

Für dieses Beispiel fügen wir hinzu Toree Kernel zur bestehenden Jupyter-Installation.

Installation von Spark + Toree + Jupyter

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

Wie aus dem angehängten Notebook hervorgeht, wird das UDF sowohl in Spark als auch in PySpark verwendet. Wir hoffen, dass dieser Teil für diejenigen von Nutzen sein wird, die Scala mögen und Modelle für maschinelles Lernen in der Produktion bereitstellen möchten.

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

Die nächsten Schritte

Obwohl sich MLflow zum Zeitpunkt der Erstellung dieses Artikels in der Alpha-Phase befindet, zeigt es vielversprechende Ansätze. Allein die Möglichkeit, mehrere Machine-Learning-Frameworks zu starten und diese über einen einzigen Endpunkt zu verwenden, hebt Empfehlungssysteme auf ein neues Level.

Darüber hinaus bringt MLflow Data Engineers und Data Scientists zusammen, indem es eine gemeinsame Schicht zwischen ihnen schafft.

Nach dieser Untersuchung von MLflow sind wir überzeugt, dass wir weitergehen und es für unsere Spark-Pipelines und in Empfehlungssystemen einsetzen werden.

Es wäre hilfreich, den Dateispeicher mit der Datenbank statt mit dem Dateisystem zu synchronisieren. So sollten wir mehrere Endpunkte erhalten, die denselben Dateispeicher verwenden können. Zum Beispiel mehrere Instanzen von Presto und Athena mit demselben Glue-Metastore.

Zusammenfassend möchten wir dem MLFlow-Community danken, dass sie unsere Arbeit mit Daten spannender macht.

Wenn Sie mit MLflow experimentieren, zögern Sie nicht, uns zu schreiben und uns zu erzählen, wie Sie es nutzen, und insbesondere, wenn Sie es in der Produktion verwenden.

Erfahren Sie mehr über die Kurse:
Machine Learning. Grundkurs
Machine Learning. Fortgeschrittener Kurs

Weiterlesen:

Quelle: habr.com

Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen 🔥 Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen | ProHoster