Здравейте, хабровчани. Както вече писахме, този месец OTUS стартира веднага два курса по машинно обучение, а именно и . В тази връзка продължаваме да споделяме полезен материал.
Целта на тази статия е да разкаже за нашия първи опит с използването на .
Ще започнем преглед с неговия tracking-сървър и ще логваме всички итерации на изследването. След това ще споделим опита си от свързването на Spark с MLflow посредством UDF.
Контекст
Ние в използваме машинно обучение и изкуствен интелект, за да дадем възможност на хората да се грижат за своето здраве и благополучие. Затова моделите на машинното обучение са в основата на нашите разработвани продукти за обработка на данни, и именно затова привлече вниманието ни MLflow — платформа с отворен код, която обхваща всички аспекти на жизнения цикъл на машинното обучение.
MLflow
Основната цел на MLflow е да осигури допълнителен слой над машинното обучение, който да позволи на специалистите по data science да работят практически с всяка библиотека за машинно обучение (, , , , и ), издигайки нейната работа на ново ниво.
MLflow предлага три компонента:
- Tracking – запис и заявки за експерименти: код, данни, конфигурация и резултати. Следенето на процеса на създаване на модел е изключително важно.
- Проекти – Формат на опаковане за стартиране на всяка платформа (например, )
- Models – общ формат за изпращане на модели в различни инструменти за разполагане.
MLflow (в момента на написване на статията в alpha версия) е платформа с отворен код, която позволява управление на жизнения цикъл на машинното обучение, включително експерименти, повторна употреба и разполагане.
Настройване на MLflow
За да използвате MLflow, първо трябва да настроите цялата среда на Python, за което ще използваме (за да инсталирате Python на Mac, вижте ). Така можем да създадем виртуална среда, където ще инсталираме всички необходими библиотеки за стартиране.
```
pyenv install 3.7.0
pyenv global 3.7.0 # Използвайте Python 3.7
mkvirtualenv mlflow # Създайте виртуална среда с Python 3.7
workon mlflow
```Ще инсталираме необходимите библиотеки.
```
pip install mlflow==0.7.0
Cython==0.29
numpy==1.14.5
pandas==0.23.4
pyarrow==0.11.0
```Забележка: използваме PyArrow за стартиране на модели като UDF. Версиите на PyArrow и Numpy трябваше да се коригират, тъй като последните версии са в конфликт помежду си.
Стартираме Tracking UI
MLflow Tracking ни позволява да водим записи и да извършваме заявки за експерименти с помощта на Python и API. Освен това, можем да определим къде да съхраняваме артефактите на модела (localhost, , , или ). Тъй като в Alpha Health използваме AWS, артефактите ще се съхраняват в S3.
# Running a Tracking Server
mlflow server
--file-store /tmp/mlflow/fileStore
--default-artifact-root s3://<bucket>/mlflow/artifacts/
--host localhost
--port 5000 MLflow препоръчва използването на постоянно файлово хранилище. Файловото хранилище е мястото, където сървърът ще съхранява метаданните за стартирания и експериментите. При стартиране на сървъра, уверете се, че той сочи към постоянно файлово хранилище. Тук за експеримента просто ще използваме /tmp.
Имайте предвид, че ако искаме да използваме сървъра mlflow за стартиране на стари експерименти, те трябва да присъстват в файловото хранилище. Но дори и без това, можем да ги използваме в UDF, тъй като ни е нужен само пътят до модела.
Забележка: Имайте предвид, че Tracking UI и клиентът на модела трябва да имат достъп до местоположението на артефакта. Тоест независимо от факта, че Tracking UI се намира в инстанция EC2, при локално стартиране на MLflow, машината трябва да има директен достъп до S3, за да записва модели на артефакти.

Tracking UI съхранява артефактите в S3 бакет
Стартиране на модели
След като Tracking сървърът е активен, можем да започнем обучение на модели.
За пример ще използваме модификацията на wine от примера на MLflow в .
MLFLOW_TRACKING_URI=http://localhost:5000 python wine_quality.py
--alpha 0.9
--l1_ratio 0.5
--wine_file ./data/winequality-red.csvКакто вече споменахме, MLflow позволява да водим записи на параметри, метрики и артефакти на модели, за да можем да проследяваме как те се развиват с течението на итерациите. Тази функция е изключително полезна, тъй като така можем да възпроизвеждаме най-добрия модел, като се обърнем към Tracking сървъра или разберем кой код е изпълнил нужната итерация, използвайки логовете на git hash коммитите.
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") 
Итерации wine
Сървърната част за модела
Сървър за проследяване MLflow, стартиран с команда „mlflow server“, разполага с REST API за проследяване на стартирания и записване на данни в локалната файлова система. Можете да зададете адреса на сървъра за проследяване чрез променлива на средата „MLFLOW_TRACKING_URI“ и API за проследяване MLflow автоматично ще се свърже с него на този адрес, за да създаде/получи информация за стартиране, метрики, логове и т.н.
Източник:
За да осигурим модел на сървъра, ще ни е нужен стартиран сървър за проследяване (вижте интерфейса за стартиране) и Run ID на модела.

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 За да обслужваме модели с помощта на функционалността MLflow serve, ще ни е нужен достъп до интерфейса за проследяване, за да получим информация за модела, просто посочвайки --run_id.
След като моделът се свърже с сървъра за проследяване, можем да получим нова крайна точка на модела.
# 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]}Стартиране на модели от Spark
Въпреки че сървърът за проследяване е достатъчно мощен за обслужване на модели в реално време, тяхното обучение и използване на функционалността serve (източник: ), използването на Spark (пакетно или потоково) е още по-мощно решение заради разпределеността.
Представете си, че просто сте провели обучение офлайн и след това приложили изходния модел върху всички ваши данни. Тук точно Spark и MLflow ще покажат най-доброто от себе си.
Инсталираме PySpark + Jupyter + Spark
Източник:
За да покажем как прилагаме модели MLflow към Spark DataFrames, трябва да настроим съвместната работа на Jupyter notebooks с PySpark.
Започнете с инсталирането на последната стабилна версия :
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Инсталирайте PySpark и Jupyter във виртуална среда:
pip install pyspark jupyterНастройте променливи на средата:
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" Определяйки notebook-dir, можем да съхраняваме нашите notebooks в желаната папка.
Стартираме Jupyter от PySpark
Тъй като успяхме да настроим Jupyter като драйвер на PySpark, сега можем да стартираме Jupyter notebook в контекста на PySpark.
(mlflow) afranzi:~$ pyspark
[I 19:05:01.572 NotebookApp] Включено расширение sparkmagic!
[I 19:05:01.573 NotebookApp] Ноутбуки обслуживаются из локального каталога: /Users/afranzi/Projects/notebooks
[I 19:05:01.573 NotebookApp] Jupyter Notebook запущен по адресу:
[I 19:05:01.573 NotebookApp] http://localhost:8888/?token=c06252daa6a12cfdd33c1d2e96c8d3b19d90e9f6fc171745
[I 19:05:01.573 NotebookApp] Используйте Control-C, чтобы остановить этот сервер и завершить все ядра (дважды, чтобы пропустить подтверждение).
[C 19:05:01.574 NotebookApp]
Скопируйте/вставьте этот URL в ваш браузер при первом подключении,
чтобы войти с токеном:
http://localhost:8888/?token=c06252daa6a12cfdd33c1d2e96c8d3b19d90e9f6fc171745 
Как было упомянуто ранее, MLflow предлагает функцию логирования артефактов модели в S3. Как только у нас есть выбранная модель, мы можем импортировать её как UDF с помощью модуля 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) 
PySpark – Прогноз качества вина
На данный момент мы говорили о том, как использовать PySpark с MLflow для предсказания качества вина на всем наборе данных wine. Но что делать, если нужно использовать модули Python MLflow из Scala Spark?
Мы протестировали это, разделив контекст Spark между Scala и Python. То есть, мы зарегистрировали MLflow UDF в Python и использовали его из Scala (да, возможно, это не лучшее решение, но что имеем).
Scala Spark + MLflow
Для этого примера мы добавим в существующий Jupyter.
Устанавливаем Spark + Toree + Jupyter
pip install toree
jupyter toree install --spark_home=${SPARK_HOME} --sys-prefix
jupyter kernelspec list
```
```
Доступные ядра:
apache_toree_scala /Users/afranzi/.virtualenvs/mlflow/share/jupyter/kernels/apache_toree_scala
python3 /Users/afranzi/.virtualenvs/mlflow/share/jupyter/kernels/python3
```Как видно из приложенного ноутбука, UDF используется как в Spark, так и в PySpark. Мы надеемся, что эта часть будет полезна тем, кто любит Scala и хочет развернуть модели машинного обучения на продакшене.
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 |
+-----------+--------+-----------+---------+-----------+
Следващи стъпки
Въпреки че в момента на написване на статията MLflow е в Alpha версия, тя изглежда доста обещаваща. Само възможността да стартирате няколко фреймворка за машинно обучение и да ги използвате от една крайна точка повишава нивото на системите за препоръки.
Освен това, MLflow сближава инженерите по данни и специалистите по данни, като прокарва общ слой между тях.
След това проучване на MLflow, сме уверени, че ще продължим напред и ще го използваме за нашите Spark потоци и в системи за препоръки.
Би било добре да синхронизираме хранилището с базата данни, вместо с файловата система. По този начин можем да получим няколко крайни точки, които да използват едно и също хранилище. Например, да използваме няколко екземпляра и с един и същ Glue metastore.
В заключение, искаме да благодарим на общността MLFlow за това, че правите работата ни с данни по-интересна.
Ако експериментирате с MLflow, не се колебайте да ни пишете и да ни разкажете как го използвате, особено ако го ползвате в продукция.
Научете повече за курсовете:
Прочетете още:
Източник: habr.com
