Erweiterung der Möglichkeiten von Spark mit MLflow

Hallo, Habr-Nutzer. Wie bereits erwähnt, starten wir in diesem Monat bei OTUS gleich zwei Kurse im Bereich Maschinelles Lernen, nämlich grundlegend und fortgeschritten. In diesem Zusammenhang setzen wir unsere regelmäßigen Informationen mit nützlichen Materialien fort.

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

zu berichten. Wir beginnen mit der Übersicht über dessen Tracking-Server und protokollieren alle Iterationen der Forschung. Dann teilen wir unsere Erfahrungen mit der Verbindung von Spark und MLflow über UDF. MLflow Alpha Health

Kontext

Wir bei setzt Maschinelles Lernen und Künstliche Intelligenz ein, um den Menschen zu helfen, sich um ihre Gesundheit und ihr Wohlbefinden zu kümmern. Daher sind die Modelle für Maschinelles Lernen die Grundlage unserer entwickelten Datenverarbeitungsprodukte, und genau deshalb hat unsere Aufmerksamkeit MLflow erregt – eine Open-Source-Plattform, die alle Aspekte des Lebenszyklus des Maschinellen Lernens abdeckt. Das Hauptziel von MLflow ist es, eine zusätzliche Schicht über dem Maschinellen Lernen bereitzustellen, die es Data-Sci entisten ermöglicht, praktisch mit jeder Bibliothek für Maschinelles Lernen zu arbeiten (

MLflow

h2omleap, keras, pytorch, tensorflow, sklearn und ), und deren Funktionalität auf ein neues Niveau zu heben.MLflow bietet drei Komponenten:

Tracking

  • – Aufzeichnung und Abfragen von Experimenten: Code, Daten, Konfiguration und Ergebnisse. Den Prozess der Modellerstellung zu verfolgen, ist sehr wichtig. – Ein Verpackungsformat, das auf jeder Plattform (z. B.
  • Projects SageMaker Models)
  • – ein allgemeines Format zum Bereitstellen von Modellen in verschiedenen Bereitstellungstools. MLflow (zum Zeitpunkt der Erstellung des Artikels in der Alpha-Version) ist eine Open-Source-Plattform, die das Management des Lebenszyklus des Maschinellen Lernens ermöglicht, einschließlich Experimenten, Wiederverwendung und Bereitstellung.

Konfiguration von MLflow

Um MLflow zu nutzen, müssen wir zunächst die gesamte Python-Umgebung konfigurieren. Dafür verwenden wir

PyEnv (um Python auf Mac zu installieren, schauen Sie bitte ). So können wir eine virtuelle Umgebung erstellen, in der wir alle benötigten Bibliotheken installieren. hier``` pyenv install 3.7.0 pyenv global 3.7.0 # Verwende Python 3.7 mkvirtualenv mlflow # Erstelle eine virtuelle Umgebung mit Python 3.7 workon mlflow ```

Jetzt installieren wir die benötigten 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 standen.

Starten wir die Tracking UI.

Starten Sie die Tracking-Benutzeroberfläche

MLflow Tracking ermöglicht es uns, Experimente mit Python zu protokollieren und Abfragen zu stellen über die 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 verwenden, wird S3 als Speicherort für die Artefakte dienen.

# 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 Dateispeichers. Ein Dateispeicher ist ein Ort, an dem der Server Metadaten zu Ausführungen und Experimenten speichert. Stellen Sie beim Starten des Servers sicher, dass er auf einen permanenten Dateispeicher verweist. Für das Experiment verwenden wir hier einfach /tmp.

Bitte beachten Sie, dass, wenn wir den MLflow-Server zum Ausführen alter Experimente verwenden möchten, diese im Dateispeicher vorhanden sein müssen. Aber auch ohne das könnten wir sie in UDF verwenden, da wir nur den Pfad zum Modell benötigen.

Hinweis: Beachten Sie, dass die Tracking UI und der Modellclient Zugriff auf den Speicherort des Artefakts haben müssen. Das bedeutet, dass unabhängig davon, dass die Tracking UI sich in einer EC2-Instanz befindet, der Computer beim lokalen Start von MLflow direkten Zugriff auf S3 haben muss, um Modelle zu schreiben.

Erweiterung der Möglichkeiten von Spark mit MLflow
Die Tracking UI speichert Artefakte in einem S3-Bucket

Modelle ausführen

Sobald der Tracking-Server läuft, können wir mit dem Training der Modelle beginnen.

Als Beispiel verwenden wir die Modifikation wine aus dem MLflow-Beispiel 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 Modellartefakten, um zu verfolgen, wie sie sich über die Iterationen entwickelt haben. Diese Funktion ist äußerst nützlich, da wir so das beste Modell reproduzieren können, indem wir auf den Tracking-Server zugreifen oder verstehen, welcher Code die gewünschte Iteration ausgeführt hat, indem wir die Protokolle der Git-Commits verwenden.

with 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', 'wine')
    mlflow.set_tag('predict', 'quality')
    mlflow.sklearn.log_model(lr, "model")

Erweiterung der Möglichkeiten von Spark mit MLflow
Iterationen wine

Backend für das Modell

Der Tracking-Server MLflow, der mit dem Befehl „mlflow server“ gestartet wurde, verfügt über eine REST-API zur Verfolgung von Ausführungen und zum Speichern 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 über diese Adresse mit dem Tracking-Server, um Informationen über Ausführungen, Metriken, Logs usw. zu erstellen/abzurufen.

Quelle: Dokumente// Einen Tracking-Server ausführen

Um ein Modell-Server bereitzustellen, benötigen wir einen laufenden Tracking-Server (siehe Schnittstelle zum Starten) und die Run ID des Modells.

Erweiterung der Möglichkeiten von Spark 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

Um Modelle mit der Funktionalität von MLflow serve zu bedienen, benötigen wir Zugriff auf die Tracking UI, um Informationen über das 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]}

Modelle aus Spark ausführen

Obwohl der Tracking-Server leistungsstark genug ist, um Modelle in Echtzeit zu bedienen, ist das Training und die Nutzung der Funktionalität serve (Quelle: mlflow // docs // models # local), ist die Nutzung von Spark (Batch oder Streaming) eine noch leistungsstärkere Lösung aufgrund der Verteilung.

Stellen Sie sich vor, Sie haben gerade ein Training offline durchgeführt und angewendet. Jetzt wenden Sie das Ausgabe-Modell auf alle Ihre Daten an. Hier zeigen Spark und MLflow ihre Stärken.

PySpark + Jupyter + Spark installieren

Quelle: Erste Schritte mit PySpark — Jupyter

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

Beginnen Sie mit der Installation der letzten 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

Umgebungsvariablen einrichten:

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 Festlegung von notebook-dir, können wir unsere Notebooks im gewünschten Ordner speichern.

Jupyter aus PySpark starten

Da wir Jupyter als PySpark-Treiber eingerichtet haben, können wir jetzt Jupyter-Notebooks 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/Einfügen Sie diese URL in Ihren Browser, wenn Sie sich das erste Mal verbinden,
    um sich mit einem Token anzumelden:
        http://localhost:8888/?token=c06252daa6a12cfdd33c1d2e96c8d3b19d90e9f6fc171745

Erweiterung der Möglichkeiten von Spark mit MLflow

Wie oben erwähnt, bietet MLflow eine Funktion zum Protokollieren von Modellartefakten in S3. Sobald wir das ausgewählte Modell in den Händen halten, haben wir die Möglichkeit, es als UDF 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 Möglichkeiten von Spark mit MLflow
PySpark – Ausgabe der Weinqualität-Prognose

Bis zu diesem Punkt haben wir darüber gesprochen, wie man PySpark mit MLflow verwendet, um die Weinqualitätsprognose auf dem gesamten Weindatensatz auszuführen. Aber was ist, wenn man die MLflow-Module in Scala Spark verwenden muss?

Wir haben auch dies getestet, indem wir den Spark-Kontext zwischen Scala und Python geteilt haben. Das heißt, wir haben die MLflow UDF in Python registriert und sie aus Scala verwendet (ja, möglicherweise keine beste Lösung, aber was haben wir).

Scala Spark + MLflow

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

Installieren von Spark + Toree + Jupyter

pip install toree
jupyter toree install --spark_home=${SPARK_HOME} --sys-prefix
jupyter kernelspec list
```
```
Verfügbare Kerne:
  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 zu erkennen ist, wird UDF sowohl mit Spark als auch mit PySpark verwendet. Wir hoffen, dass dieser Teil für diejenigen von Nutzen ist, die Scala lieben und Modelle für maschinelles Lernen in der Produktion implementieren 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       |
+-----------+--------+-----------+---------+-----------+

Nächste Schritte

Obwohl sich MLflow zum Zeitpunkt des Schreibens des Artikels in der Alpha-Version befindet, sieht es vielversprechend aus. Allein die Möglichkeit, mehrere Frameworks für maschinelles Lernen zu verwenden und sie von einem einzigen Endpunkt aus zu nutzen, hebt Empfehlungssysteme auf ein neues Niveau.

Darüber hinaus bringt MLflow Dateningenieure und Data-Science-Experten zusammen, indem es eine gemeinsame Schicht zwischen ihnen schafft.

Nach dieser MLflow-Studie sind wir zuversichtlich, dass wir weiterarbeiten werden und es für unsere Spark-Pipelines und in Empfehlungssystemen nutzen werden.

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

Zusammenfassend möchten wir dem MLflow-Team danken, dass sie unsere Arbeit mit Daten interessanter gestalten.

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

Mehr über die Kurse erfahren:
Machine Learning. Grundkurs
Machine Learning. Fortgeschrittener Kurs

Weiterlesen:

Quelle: habr.com

60GB SSD 8Gb DDR4