Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
50 commits
Select commit Hold shift + click to select a range
46c45ea
Selectable code env for MLFlow models
cstenac Sep 22, 2021
e76cc90
Class to interact with MLFlow model metrics params
cstenac Sep 23, 2021
a8243cc
Specify the name when creating a SM for an MLFlow model, rather than …
lpenet Sep 28, 2021
7306afc
Clean mlflow tmp file
Oct 7, 2021
45cde3a
Added FMClient and get_cloud_credentials
Apr 16, 2021
287fc31
Set tenant_id in FMClient
May 27, 2021
be994f2
Added FMVirtualNetwork
May 27, 2021
f7c241a
Added FMInstanceSettingsTemplate
May 27, 2021
0fda43b
Added FMInstance
May 27, 2021
7dcf67a
fm public api: Add instance actions
Jun 1, 2021
e6210ad
fm: Add FMInstanceStatus
Jun 2, 2021
a53fbfe
fm: Save and restart_dss
Jun 2, 2021
b49c8e3
fm: create instance template
Jun 2, 2021
9825ca7
FM: future
Jun 3, 2021
d1ee8d6
fm: delete instance
Jun 3, 2021
a221e17
fm: Add delete instance settings template
Jun 3, 2021
2e62ecf
fm: Virtual network management
Jun 3, 2021
4627c9c
fm: Instance settings update
Jun 3, 2021
dfe7a85
Add __init__.py in fm module
Sep 1, 2021
3a484b5
FM Tenant: Update license
Sep 1, 2021
0fbbf2e
FM: Add helper to set cloud credentials
Sep 1, 2021
412c2d0
FM: Review comments
Sep 1, 2021
0ba05b7
Add run_ansible_task & install_system_packages
Sep 3, 2021
791e53e
FM: Setup Advanced Security SetupAction
Sep 3, 2021
381030f
FM: Install JDBC Driver SetupAction
Sep 3, 2021
3e93db8
FM: setup_k8s_and_spark SetupAction
Sep 3, 2021
66a7a9e
FM: FMInstanceCreator
Sep 7, 2021
2f09e13
FM: Split Instance and InstanceCreator per cloud
Sep 8, 2021
68e4005
FM InstanceSettingsTemplateCreator
Sep 17, 2021
e257c5b
FM: InstanceSettingsTemplate format + doc
Sep 21, 2021
16c2325
fm: FMVirtualNetworkCreator
Sep 21, 2021
0611d8d
FM: format
Sep 21, 2021
26185d1
FM: Split client by cloud
Sep 21, 2021
6182ba8
FM: Update Cloud Tags
Sep 21, 2021
b64b171
typo
FChataigner Sep 22, 2021
e7b1f5a
fix created object
FChataigner Sep 22, 2021
0e9e58c
type returned objects to the appropriate subclass
FChataigner Sep 22, 2021
a766032
fix client passed to created object
FChataigner Sep 22, 2021
e3947dd
FM: Simplify CloudTags
Sep 23, 2021
9e2b9bb
FM: Review feedback on VirtualNetworks
Sep 23, 2021
a8cf85a
FM: Review feedback on InstanceSettingsTemplates
Sep 23, 2021
fb41a80
FM: Review feedback on Instances
Sep 23, 2021
59fb688
Fix import Enum
Sep 28, 2021
603f628
Add FM client to packages
Sep 28, 2021
2abe499
public api for the list of images
FChataigner Sep 28, 2021
6c901de
let VN be created with the default values
FChataigner Sep 28, 2021
21e1abd
expose actions on instances
FChataigner Sep 28, 2021
509c2fc
expose snapshots
FChataigner Sep 28, 2021
8646164
Rename trainDiagnostics into mlDiagnostics
Oct 8, 2021
4e7d77e
Trading Run IDs for Evaluation IDs (#179)
lpenet Oct 11, 2021
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion dataikuapi/__init__.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
from .dssclient import DSSClient
from .fmclient import FMClientAWS, FMClientAzure

from .apinode_client import APINodeClient
from .apinode_admin_client import APINodeAdminClient

from .dss.recipe import GroupingRecipeCreator, JoinRecipeCreator, StackRecipeCreator, WindowRecipeCreator, SyncRecipeCreator, SamplingRecipeCreator, SQLQueryRecipeCreator, CodeRecipeCreator, SplitRecipeCreator, SortRecipeCreator, TopNRecipeCreator, DistinctRecipeCreator, DownloadRecipeCreator, PredictionScoringRecipeCreator, ClusteringScoringRecipeCreator

from .dss.admin import DSSUserImpersonationRule, DSSGroupImpersonationRule
from .dss.admin import DSSUserImpersonationRule, DSSGroupImpersonationRule
2 changes: 1 addition & 1 deletion dataikuapi/dss/ml.py
Original file line number Diff line number Diff line change
Expand Up @@ -1867,7 +1867,7 @@ def get_diagnostics(self):
:returns: list of diagnostics
:rtype: list of type `dataikuapi.dss.ml.DSSMLDiagnostic`
"""
diagnostics = self.details.get("trainDiagnostics", {})
diagnostics = self.details.get("mlDiagnostics", {})
return [DSSMLDiagnostic(d) for d in diagnostics.get("diagnostics", [])]

def generate_documentation(self, folder_id=None, path=None):
Expand Down
46 changes: 23 additions & 23 deletions dataikuapi/dss/modelevaluationstore.py
Original file line number Diff line number Diff line change
Expand Up @@ -124,18 +124,18 @@ def list_model_evaluations(self):
:returns: The list of the model evaluations
:rtype: list of :class:`dataikuapi.dss.modelevaluationstore.DSSModelEvaluation`
"""
items = self.client._perform_json("GET", "/projects/%s/modelevaluationstores/%s/runs/" % (self.project_key, self.mes_id))
return [DSSModelEvaluation(self, item["ref"]["runId"]) for item in items]
items = self.client._perform_json("GET", "/projects/%s/modelevaluationstores/%s/evaluations/" % (self.project_key, self.mes_id))
return [DSSModelEvaluation(self, item["ref"]["evaluationId"]) for item in items]

def get_model_evaluation(self, run_id):
def get_model_evaluation(self, evaluation_id):
"""
Get a handle to interact with a specific model evaluation

:param string run_id: the id of the desired model evaluation
:param string evaluation_id: the id of the desired model evaluation

:returns: A :class:`dataikuapi.dss.modelevaluationstore.DSSModelEvaluation` model evaluation handle
"""
return DSSModelEvaluation(self, run_id)
return DSSModelEvaluation(self, evaluation_id)

def get_latest_model_evaluation(self):
"""
Expand All @@ -146,11 +146,11 @@ def get_latest_model_evaluation(self):
if the store is not empty, else None
"""

latest_run_id = self.client._perform_text(
"GET", "/projects/%s/modelevaluationstores/%s/latestRunId" % (self.project_key, self.mes_id))
if not latest_run_id:
latest_evaluation_id = self.client._perform_text(
"GET", "/projects/%s/modelevaluationstores/%s/latestEvaluationId" % (self.project_key, self.mes_id))
if not latest_evaluation_id:
return None
return DSSModelEvaluation(self, latest_run_id)
return DSSModelEvaluation(self, latest_evaluation_id)

def delete_model_evaluations(self, evaluations):
"""
Expand All @@ -159,13 +159,13 @@ def delete_model_evaluations(self, evaluations):
obj = []
for evaluation in evaluations:
if isinstance(evaluation, DSSModelEvaluation):
obj.append(evaluation.run_id)
obj.append(evaluation.evaluation_id)
elif isinstance(evaluation, dict):
obj.append(evaluation['run_id'])
obj.append(evaluation['evaluation_id'])
else:
obj.append(evaluation)
self.client._perform_json(
"DELETE", "/projects/%s/modelevaluationstores/%s/runs/" % (self.project_key, self.mes_id, self.run_id), body=obj)
"DELETE", "/projects/%s/modelevaluationstores/%s/evaluations/" % (self.project_key, self.mes_id), body=obj)

def build(self, job_type="NON_RECURSIVE_FORCED_BUILD", wait=True, no_fail=False):
"""
Expand Down Expand Up @@ -263,11 +263,11 @@ class DSSModelEvaluation:
Do not create this class directly, instead use :meth:`dataikuapi.dss.DSSModelEvaluationStore.get_model_evaluation`
"""

def __init__(self, model_evaluation_store, run_id):
def __init__(self, model_evaluation_store, evaluation_id):
self.model_evaluation_store = model_evaluation_store
self.client = model_evaluation_store.client
# unpack some fields
self.run_id = run_id
self.evaluation_id = evaluation_id
self.project_key = model_evaluation_store.project_key
self.mes_id = model_evaluation_store.mes_id

Expand All @@ -276,23 +276,23 @@ def get_full_info(self):
Retrieve the model evaluation with its performance data
"""
data = self.client._perform_json(
"GET", "/projects/%s/modelevaluationstores/%s/runs/%s" % (self.project_key, self.mes_id, self.run_id))
"GET", "/projects/%s/modelevaluationstores/%s/evaluations/%s" % (self.project_key, self.mes_id, self.evaluation_id))
return DSSModelEvaluationFullInfo(self, data)

def get_full_id(self):
return "ME-{}-{}-{}".format(self.project_key, self.mes_id, self.run_id)
return "ME-{}-{}-{}".format(self.project_key, self.mes_id, self.evaluation_id)

def delete(self):
"""
Remove this model evaluation
"""
obj = [self.run_id]
obj = [self.evaluation_id]
self.client._perform_json(
"DELETE", "/projects/%s/modelevaluationstores/%s/runs/" % (self.project_key, self.mes_id), body=obj)
"DELETE", "/projects/%s/modelevaluationstores/%s/evaluations/" % (self.project_key, self.mes_id), body=obj)

@property
def full_id(self):
return "ME-%s-%s-%s"%(self.project_key, self.mes_id, self.run_id)
return "ME-%s-%s-%s"%(self.project_key, self.mes_id, self.evaluation_id)

def compute_data_drift(self, reference=None, data_drift_params=None, wait=True):
"""
Expand All @@ -310,7 +310,7 @@ def compute_data_drift(self, reference=None, data_drift_params=None, wait=True):
reference = reference.full_id

future_response = self.client._perform_json(
"POST", "/projects/%s/modelevaluationstores/%s/runs/%s/computeDataDrift" % (self.project_key, self.mes_id, self.run_id),
"POST", "/projects/%s/modelevaluationstores/%s/evaluations/%s/computeDataDrift" % (self.project_key, self.mes_id, self.evaluation_id),
body={
"referenceId": reference,
"dataDriftParams": data_drift_params
Expand All @@ -325,7 +325,7 @@ def get_metrics(self):
:return: the metrics, as a JSON object
"""
return self.client._perform_json(
"GET", "/projects/%s/modelevaluationstores/%s/runs/%s/metrics" % (self.project_key, self.mes_id, self.run_id))
"GET", "/projects/%s/modelevaluationstores/%s/evaluations/%s/metrics" % (self.project_key, self.mes_id, self.evaluation_id))

def get_sample_df(self):
"""
Expand All @@ -337,12 +337,12 @@ def get_sample_df(self):
buf = BytesIO()
with self.client._perform_raw(
"GET",
"/projects/%s/modelevaluationstores/%s/runs/%s/sample" % (self.project_key, self.mes_id, self.run_id)
"/projects/%s/modelevaluationstores/%s/evaluations/%s/sample" % (self.project_key, self.mes_id, self.evaluation_id)
).raw as f:
buf.write(f.read())
schema_txt = self.client._perform_raw(
"GET",
"/projects/%s/modelevaluationstores/%s/runs/%s/schema" % (self.project_key, self.mes_id, self.run_id)
"/projects/%s/modelevaluationstores/%s/evaluations/%s/schema" % (self.project_key, self.mes_id, self.evaluation_id)
).text
schema = json.loads(schema_txt)
import pandas as pd
Expand Down
14 changes: 7 additions & 7 deletions dataikuapi/dss/project.py
Original file line number Diff line number Diff line change
Expand Up @@ -709,22 +709,22 @@ def get_saved_model(self, sm_id):
"""
return DSSSavedModel(self.client, self.project_key, sm_id)

def create_mlflow_pyfunc_model(self, id, prediction_type = None):
def create_mlflow_pyfunc_model(self, name, prediction_type = None):
"""
Creates a new external saved model for storing and managing MLFlow models

:param string id: Identifier for the new saved model in the flow
:param string name: Human readable name for the new saved model in the flow
:param string prediction_type: Optional (but needed for most operations). One of BINARY_CLASSIFICATION, MULTICLASS or REGRESSION
"""
if len(id) != 8:
raise ValueError("model id must be 8 characters long")
if not name:
raise ValueError("name can not be empty")
model = {
"id": id,
"savedModelType" : "MLFLOW_PYFUNC",
"predictionType" : prediction_type
"predictionType" : prediction_type,
"name": name
}

self.client._perform_empty("POST", "/projects/%s/savedmodels/" % self.project_key, body = model)
id = self.client._perform_json("POST", "/projects/%s/savedmodels/" % self.project_key, body = model)["id"]
return self.get_saved_model(id)

########################################################
Expand Down
30 changes: 27 additions & 3 deletions dataikuapi/dss/savedmodel.py
Original file line number Diff line number Diff line change
Expand Up @@ -117,26 +117,30 @@ def get_origin_ml_task(self):
if fmi is not None:
return DSSMLTask.from_full_model_id(self.client, fmi, project_key=self.project_key)

def import_mlflow_version_from_path(self, version_id, path):
def import_mlflow_version_from_path(self, version_id, path, code_env_name = "INHERIT"):
"""
Create a new version for this saved model from a path containing a MLFlow model.

Requires the saved model to have been created using :meth:`dataikuapi.dss.project.DSSProject.create_mlflow_pyfunc_model`.

:param str version_id: Identifier of the version to create
:param str path: An absolute path on the local filesystem. Must be a folder, and must contain a MLFlow model

:param str code_env_name: Name of the code env to use for this model version. The code env must contain at least
mlflow and the package(s) corresponding to the used MLFlow-compatible frameworks.
If value is "INHERIT", the default active code env of the project will be used
:return a :class:MLFlowVersionHandler in order to interact with the new MLFlow model version
"""
# TODO: Add a check that it's indeed a MLFlow model folder
# TODO: Put it in a proper temp folder
# TODO: cleanup the archive
import shutil
import os
shutil.make_archive("tmpmodel", "zip", path) #[, root_dir[, base_dir[, verbose[, dry_run[, owner[, group[, logger]]]]]]])

with open("tmpmodel.zip", "rb") as fp:
self.client._perform_empty("POST", "/projects/%s/savedmodels/%s/versions/%s" % (self.project_key, self.sm_id, version_id),
self.client._perform_empty("POST", "/projects/%s/savedmodels/%s/versions/%s?codeEnvName=%s" % (self.project_key, self.sm_id, version_id, code_env_name),
files={"file":("tmpmodel.zip", fp)})
os.remove("tmpmodel.zip")

return self.get_mlflow_version_handler(version_id)

Expand Down Expand Up @@ -232,13 +236,33 @@ def delete(self):
"""
return self.client._perform_empty("DELETE", "/projects/%s/savedmodels/%s" % (self.project_key, self.sm_id))

class MLFlowVersionSettings:
"""Handle for the settings of an imported MLFlow model version"""

def __init__(self, version_handler, data):
self.version_handler = version_handler
self.data = data

@property
def raw(self):
return self.data

def save(self):
self.version_handler.saved_model.client._perform_empty("PUT",
"/projects/%s/savedmodels/%s/versions/%s/external-ml/metadata" % (self.version_handler.saved_model.project_key, self.version_handler.saved_model.sm_id, self.version_handler.version_id),
body=self.data)

class MLFlowVersionHandler:
"""Handler to interact with an imported MLFlow model version"""
def __init__(self, saved_model, version_id):
"""Do not call this, use :meth:`DSSSavedModel.get_mlflow_version_handler`"""
self.saved_model = saved_model
self.version_id = version_id

def get_settings(self):
metadata = self.saved_model.client._perform_json("GET", "/projects/%s/savedmodels/%s/versions/%s/external-ml/metadata" % (self.saved_model.project_key, self.saved_model.sm_id, self.version_id))
return MLFlowVersionSettings(self, metadata)

def set_core_metadata(self,
target_column_name, class_labels = None,
get_features_from_dataset=None, features_list = None,
Expand Down
Empty file added dataikuapi/fm/__init__.py
Empty file.
101 changes: 101 additions & 0 deletions dataikuapi/fm/future.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
import sys, time


class FMFuture(object):
"""
A future on the DSS instance
"""

def __init__(
self, client, job_id, state=None, result_wrapper=lambda result: result
):
self.client = client
self.job_id = job_id
self.state = state
self.state_is_peek = True
self.result_wrapper = result_wrapper

@staticmethod
def from_resp(client, resp, result_wrapper=lambda result: result):
"""Creates a DSSFuture from a parsed JSON response"""
return FMFuture(
client, resp.get("jobId", None), state=resp, result_wrapper=result_wrapper
)

@classmethod
def get_result_wait_if_needed(cls, client, ret):
if "jobId" in ret:
future = FMFuture(client, ret["jobId"], ret)
future.wait_for_result()
return future.get_result()
else:
return ret["result"]

def abort(self):
"""
Abort the future
"""
return self.client._perform_tenant_empty("DELETE", "/futures/%s" % self.job_id)

def get_state(self):
"""
Get the status of the future, and its result if it's ready
"""
self.state = self.client._perform_tenant_json(
"GET", "/futures/%s" % self.job_id, params={"peek": False}
)
self.state_is_peek = False
return self.state

def peek_state(self):
"""
Get the status of the future, and its result if it's ready
"""
self.state = self.client._perform_tenant_json(
"GET", "/futures/%s" % self.job_id, params={"peek": True}
)
self.state_is_peek = True
return self.state

def get_result(self):
"""
Get the future result if it's ready, raises an Exception otherwise
"""
if (
self.state is None
or not self.state.get("hasResult", False)
or self.state_is_peek
):
self.get_state()
if self.state.get("hasResult", False):
return self.result_wrapper(self.state.get("result", None))
else:
raise Exception("Result not ready")

def has_result(self):
"""
Checks whether the future has a result ready
"""
if self.state is None or not self.state.get("hasResult", False):
self.get_state()
return self.state.get("hasResult", False)

def wait_for_result(self):
"""
Wait and get the future result
"""
if self.state.get("hasResult", False):
return self.result_wrapper(self.state.get("result", None))
if (
self.state is None
or not self.state.get("hasResult", False)
or self.state_is_peek
):
self.get_state()
while not self.state.get("hasResult", False):
time.sleep(5)
self.get_state()
if self.state.get("hasResult", False):
return self.result_wrapper(self.state.get("result", None))
else:
raise Exception("No result")
Loading