diff --git a/MLproject b/MLproject index 0f59f56..2ec5e09 100644 --- a/MLproject +++ b/MLproject @@ -11,34 +11,26 @@ entry_points: command: "python generation.py --num_samples {num_samples} --num_frames {num_frames} --exposure_time {exposure_time}" analysis1: parameters: - generated_data: {type: string, default: "/tmp/foobar"} - num_samples: {type: int, default: 1} - num_frames: {type: int, default: 5} + generation: {type: string, default: ""} max_sigma: {type: int, default: 4} min_sigma: {type: int, default: 1} threshold: {type: float, default: 50.0} overlap: {type: float, default: 0.5} - interval: {type: float, default: 33.0e-3} - command: "python analysis1.py --generated_data {generated_data} --num_samples {num_samples} --num_frames {num_frames} --min_sigma {min_sigma} --max_sigma {max_sigma} --threshold {threshold} --overlap {overlap} --interval {interval}" + command: "python analysis1.py --generation {generation} --min_sigma {min_sigma} --max_sigma {max_sigma} --threshold {threshold} --overlap {overlap}" analysis2: parameters: - generated_data: {type: string, default: "/tmp/foobar"} - num_samples: {type: int, default: 1} - num_frames: {type: int, default: 5} - threshold: {type: float, default: 50.0} - interval: {type: float, default: 33.0e-3} - command: "python analysis2.py --generated_data {generated_data} --num_samples {num_samples} --num_frames {num_frames} --threshold {threshold} --interval {interval}" + generation: {type: string, default: ""} + analysis1: {type: string, default: ""} + max_distance: {type: float, default: 50.0} + command: "python analysis2.py --generation {generation} --analysis1 {analysis1} --max_distance {max_distance}" evaluation1: parameters: - generated_data: {type: string, default: "/tmp/foobar"} - num_samples: {type: int, default: 1} - num_frames: {type: int, default: 5} - threshold: {type: float, default: 50.0} - command: "python evaluation1.py --generated_data {generated_data} --num_samples {num_samples} --num_frames {num_frames} --threshold {threshold}" + generation: {type: string, default: ""} + analysis1: {type: string, default: ""} + command: "python evaluation1.py --generation {generation} --analysis1 {analysis1}" main: parameters: num_samples: {type: int, default: 1} num_frames: {type: int, default: 5} threshold: {type: float, default: 50.0} command: "python main.py --num_samples {num_samples} --num_frames {num_frames} --threshold {threshold}" - diff --git a/analysis1.py b/analysis1.py index 6b26c11..6f785ce 100644 --- a/analysis1.py +++ b/analysis1.py @@ -1,67 +1,56 @@ -# -*- coding: utf-8 -*- -"""analysis1.ipynb - -Automatically generated by Colaboratory. - -Original file is located at - https://colab.research.google.com/github/ecell/bioimage_workflows/blob/master/analysis1.ipynb -""" - import argparse +import pathlib + +import mlflow +from mlflow import log_metric, log_param, log_artifacts +entrypoint = "analysis1" parser = argparse.ArgumentParser(description='analysis1 step') -parser.add_argument('--generated_data', type=str, default="/tmp/foobar") -parser.add_argument('--num_samples', type=int, default=1) -parser.add_argument('--num_frames', type=int, default=5) +parser.add_argument('--generation', type=str, default="") parser.add_argument('--min_sigma', type=int, default=1) parser.add_argument('--max_sigma', type=int, default=4) parser.add_argument('--threshold', type=float, default=50.0) parser.add_argument('--overlap', type=float, default=0.5) -parser.add_argument('--interval', type=float, default=33.0e-3) - args = parser.parse_args() -import mlflow -mlflow.start_run(run_name="analysis1") +active_run = mlflow.start_run() +mlflow.set_tag("mlflow.runName", entrypoint) -generated_data = args.generated_data -num_samples = args.num_samples -num_frames = args.num_frames +generation = args.generation min_sigma = args.min_sigma max_sigma = args.max_sigma threshold = args.threshold overlap = args.overlap -interval = args.interval -from mlflow import log_metric, log_param, log_artifacts -log_param("num_frames", num_frames) -log_param("num_samples", num_samples) -log_param("min_sigma", min_sigma) -log_param("max_sigma", max_sigma) -log_param("threshold", threshold) -log_param("overlap", overlap) -log_param("interval", interval) +for key, value in vars(args).items(): + log_param(key, value) -nproc = 1 +client = mlflow.tracking.MlflowClient() +generation_run = client.get_run(generation) +num_samples = int(generation_run.data.params["num_samples"]) +num_frames = int(generation_run.data.params["num_frames"]) +interval = float(generation_run.data.params["interval"]) +generation_artifacts = pathlib.Path(client.download_artifacts(generation, ".")) + +import tempfile +artifacts = pathlib.Path(tempfile.mkdtemp()) / "artifacts" +artifacts.mkdir(parents=True, exist_ok=True) + +#XXX: HERE import numpy timepoints = numpy.linspace(0, interval * num_frames, num_frames + 1) -import pathlib -inputpath = pathlib.Path(generated_data.replace("file://", "")) -artifacts = pathlib.Path(generated_data.replace("file://", "")) -artifacts.mkdir(parents=True, exist_ok=True) - import scopyon import warnings warnings.simplefilter('ignore', RuntimeWarning) for i in range(num_samples): - imgs = [scopyon.Image(data) for data in numpy.load(inputpath / f"images{i:03d}.npy")] + imgs = [scopyon.Image(data) for data in numpy.load(generation_artifacts / f"images{i:03d}.npy")] spots = [ scopyon.analysis.spot_detection( - img.as_array(), processes=nproc, + img.as_array(), min_sigma=min_sigma, max_sigma=max_sigma, threshold=threshold, overlap=overlap) for img in imgs] @@ -70,13 +59,21 @@ spots_.extend(([t] + list(row) for row in data)) spots_ = numpy.array(spots_) numpy.save(artifacts / f"spots{i:03d}.npy", spots_) - + + r = 6 + shapes = [dict(x=spot[0], y=spot[1], sigma=r, color='red') + for spot in spots[0]] + imgs[0].save(artifacts / f"spots{i:03d}_000.png", shapes=shapes) + print("{} spots are detected in {} frames.".format(len(spots_), len(imgs))) + log_metric("num_spots", len(spots_)) warnings.resetwarnings() -#!ls ./artifacts +#XXX: THERE -#log_artifacts("./artifacts") -log_artifacts(generated_data) +log_artifacts(str(artifacts)) mlflow.end_run() + +import shutil +shutil.rmtree(str(artifacts)) diff --git a/analysis2.py b/analysis2.py index 4d7820d..e123fe7 100644 --- a/analysis2.py +++ b/analysis2.py @@ -1,52 +1,50 @@ -# -*- coding: utf-8 -*- -"""analysis2.ipynb - -Automatically generated by Colaboratory. - -Original file is located at - https://colab.research.google.com/github/ecell/bioimage_workflows/blob/master/analysis2.ipynb -""" - import argparse +import pathlib + +import mlflow +from mlflow import log_metric, log_param, log_artifacts +entrypoint = "analysis2" parser = argparse.ArgumentParser(description='analysis2 step') -parser.add_argument('--generated_data', type=str, default="/tmp/foobar") -parser.add_argument('--num_samples', type=int, default=1) -parser.add_argument('--num_frames', type=int, default=5) -parser.add_argument('--threshold', type=float, default=50.0) -parser.add_argument('--interval', type=float, default=33.0e-3) +parser.add_argument('--generation', type=str, default="") +parser.add_argument('--analysis1', type=str, default="") +parser.add_argument('--seed', type=int, default=123) +parser.add_argument('--max_distance', type=float, default=50.0) args = parser.parse_args() -import mlflow -mlflow.start_run(run_name="analysis2") +active_run = mlflow.start_run() +mlflow.set_tag("mlflow.runName", entrypoint) -generated_data = args.generated_data -num_samples = args.num_samples -num_frames = args.num_frames -threshold = args.threshold -interval = args.interval +generation = args.generation +analysis1 = args.analysis1 +seed = args.seed +max_distance = args.max_distance -seed = 123 +for key, value in vars(args).items(): + log_param(key, value) -from mlflow import log_metric, log_param, log_artifacts -log_param("num_frames", num_frames) -log_param("num_samples", num_samples) -log_param("interval", interval) -log_param("seed", seed) -log_param("threshold", threshold) +client = mlflow.tracking.MlflowClient() +generation_run = client.get_run(generation) +analysis1_run = client.get_run(analysis1) +num_samples = int(generation_run.data.params["num_samples"]) +num_frames = int(generation_run.data.params["num_frames"]) +interval = float(generation_run.data.params["interval"]) +generation_artifacts = pathlib.Path(client.download_artifacts(generation, ".")) +analysis1_artifacts = pathlib.Path(client.download_artifacts(analysis1, ".")) -import pathlib -inputpath = pathlib.Path("generated_data") -artifacts = pathlib.Path("generated_data") +import tempfile +artifacts = pathlib.Path(tempfile.mkdtemp()) / "artifacts" artifacts.mkdir(parents=True, exist_ok=True) +#XXX: HERE + import scopyon -config = scopyon.Configuration(filename=inputpath / "config.yaml") +config = scopyon.Configuration(filename=generation_artifacts / "config.yaml") pixel_length = config.default.detector.pixel_length / config.default.magnification import numpy -def trace_spots(spots, threshold=numpy.inf, ndim=2): +def trace_spots(spots, max_distance=numpy.inf, ndim=2): observation_vec = [] lengths = [] for i in range(len(spots[0])): @@ -55,7 +53,7 @@ def trace_spots(spots, threshold=numpy.inf, ndim=2): displacements = numpy.power(spots[j][:, : ndim] - spots[j - 1][iprev, : ndim], 2).sum(axis=1) inext = displacements.argmin() displacement = numpy.sqrt(displacements[inext]) - if displacement > threshold: + if displacement > max_distance: if j > 1: lengths.append(j - 1) break @@ -70,7 +68,7 @@ def trace_spots(spots, threshold=numpy.inf, ndim=2): lengths = [] ndim = 2 for i in range(num_samples): - spots_ = numpy.load(inputpath / f"spots{i:03d}.npy") + spots_ = numpy.load(analysis1_artifacts / f"spots{i:03d}.npy") t = spots_[0, 0] spots = [[spots_[0, 1: ]]] for row in spots_[1: ]: @@ -84,7 +82,7 @@ def trace_spots(spots, threshold=numpy.inf, ndim=2): spots[-1] = numpy.asarray(spots[-1]) # print(spots) - observation_vec_, lengths_ = trace_spots(spots, threshold=threshold, ndim=ndim) + observation_vec_, lengths_ = trace_spots(spots, max_distance=max_distance, ndim=ndim) observation_vec.extend(observation_vec_) lengths.extend(lengths_) observation_vec = numpy.array(observation_vec) @@ -95,20 +93,16 @@ def trace_spots(spots, threshold=numpy.inf, ndim=2): from plotly.subplots import make_subplots fig = make_subplots(rows=1, cols=2, subplot_titles=['Square Displacement', 'Intensity']) - fig.add_trace(go.Histogram(x=observation_vec[:, 0], nbinsx=30, histnorm='probability'), row=1, col=1) - fig.add_trace(go.Histogram(x=observation_vec[:, 1], nbinsx=30, histnorm='probability'), row=1, col=2) - fig.update_layout(barmode='overlay') fig.update_traces(opacity=0.75, showlegend=False) -#fig.show() -fig.write_image(generated_data + "/analysis2_1.png") +# fig.show() +fig.write_image(str(artifacts / "histogram1.png")) from scopyon.analysis import PTHMM rng = numpy.random.RandomState(seed) - model = PTHMM(n_diffusivities=3, n_oligomers=1, n_iter=100, random_state=rng) model.fit(observation_vec, lengths) @@ -132,18 +126,19 @@ def trace_spots(spots, threshold=numpy.inf, ndim=2): expected_vec[sum(lengths[: i]): sum(lengths[: i + 1])] = X_ fig = make_subplots(rows=1, cols=2, subplot_titles=['Square Displacement', 'Intensity']) - fig.add_trace(go.Histogram(x=observation_vec[:, 0], nbinsx=30, histnorm='probability density'), row=1, col=1) fig.add_trace(go.Histogram(x=expected_vec[:, 0], nbinsx=30, histnorm='probability density'), row=1, col=1) - fig.add_trace(go.Histogram(x=observation_vec[:, 1], nbinsx=30, histnorm='probability density'), row=1, col=2) fig.add_trace(go.Histogram(x=expected_vec[:, 1], nbinsx=30, histnorm='probability density'), row=1, col=2) - fig.update_layout(barmode='overlay') fig.update_traces(opacity=0.75, showlegend=False) -#fig.show() -fig.write_image(generated_data + "/analysis2_2.png") +# fig.show() +fig.write_image(str(artifacts / "histogram2.png")) -#log_artifacts("./artifacts") -log_artifacts(generated_data) +#XXX: THERE + +log_artifacts(str(artifacts)) mlflow.end_run() + +import shutil +shutil.rmtree(str(artifacts)) diff --git a/evaluation1.py b/evaluation1.py index 44a6395..d668046 100644 --- a/evaluation1.py +++ b/evaluation1.py @@ -1,52 +1,54 @@ -# -*- coding: utf-8 -*- -"""evaluation1.ipynb +import argparse +import pathlib -Automatically generated by Colaboratory. +import mlflow +from mlflow import log_metric, log_param, log_artifacts -Original file is located at - https://colab.research.google.com/github/ecell/bioimage_workflows/blob/master/evaluation1.ipynb -""" +entrypoint = "evaluation1" +parser = argparse.ArgumentParser(description='evaluation1 step') +parser.add_argument('--generation', type=str, default="") +parser.add_argument('--analysis1', type=str, default="") +# parser.add_argument('--analysis2', type=str, default="") +parser.add_argument('--max_distance', type=float, default=50.0) +args = parser.parse_args() -import argparse +active_run = mlflow.start_run() +mlflow.set_tag("mlflow.runName", entrypoint) -parser = argparse.ArgumentParser(description='evaluation1 step') -parser.add_argument('--generated_data', type=str, default="/tmp/foobar") -parser.add_argument('--num_samples', type=int, default=1) -parser.add_argument('--num_frames', type=int, default=5) -parser.add_argument('--threshold', type=float, default=50.0) +generation = args.generation +analysis1 = args.analysis1 +# analysis2 = args.analysis2 +max_distance = args.max_distance -args = parser.parse_args() +for key, value in vars(args).items(): + log_param(key, value) -import mlflow -mlflow.start_run(run_name="evaluation1") +client = mlflow.tracking.MlflowClient() +generation_run = client.get_run(generation) +num_samples = int(generation_run.data.params["num_samples"]) +analysis1_run = client.get_run(analysis1) +# analysis2_run = client.get_run(analysis2) +generation_artifacts = pathlib.Path(client.download_artifacts(generation, ".")) +analysis1_artifacts = pathlib.Path(client.download_artifacts(analysis1, ".")) +# analysis2_artifacts = pathlib.Path(client.download_artifacts(analysis2, ".")) -generated_data = args.generated_data -num_samples = args.num_samples -num_frames = args.num_frames -threshold = args.threshold +import tempfile +artifacts = pathlib.Path(tempfile.mkdtemp()) / "artifacts" +artifacts.mkdir(parents=True, exist_ok=True) -from mlflow import log_metric, log_param, log_artifacts -log_param("generated_data", generated_data) -log_param("num_samples", num_samples) -log_param("num_frames", num_frames) -log_param("threshold", threshold) +#XXX: HERE import numpy -import pathlib -inputpath = pathlib.Path(generated_data) -# artifacts = pathlib.Path("./artifacts") -# artifacts.mkdir(parents=True, exist_ok=True) - import scopyon -config = scopyon.Configuration(filename=inputpath / "config.yaml") +config = scopyon.Configuration(filename=generation_artifacts / "config.yaml") pixel_length = config.default.detector.pixel_length / config.default.magnification rates = numpy.zeros(4, dtype=int) closest = [] for i in range(num_samples): - true_data_ = numpy.load(inputpath / f"true_data{i:03d}.npy") + true_data_ = numpy.load(generation_artifacts / f"true_data{i:03d}.npy") t = true_data_[0, 0] true_data = [[true_data_[0, 1: ]]] for row in true_data_[1: ]: @@ -59,7 +61,7 @@ else: true_data[-1] = numpy.asarray(true_data[-1]) - spots_ = numpy.load(inputpath / f"spots{i:03d}.npy") + spots_ = numpy.load(analysis1_artifacts / f"spots{i:03d}.npy") t = spots_[0, 0] spots = [[spots_[0, 1: ]]] for row in spots_[1: ]: @@ -78,9 +80,9 @@ distance = data - spot[0: 2] idx = (distance ** 2).sum(axis=1).argmin() closest.append(distance[idx]) - + distance = numpy.sqrt(distance[idx] ** 2).sum() - if distance < threshold: + if distance < max_distance: rates[0] += 1 else: rates[1] += 1 @@ -90,7 +92,7 @@ distance = (distance ** 2).sum(axis=1) idx = distance.argmin() distance = numpy.sqrt(distance[idx]) - if distance < threshold: + if distance < max_distance: rates[2] += 1 else: rates[3] += 1 @@ -107,20 +109,20 @@ log_metric("x_std", x_std) log_metric("y_std", y_std) -import plotly.express as px -w = h = 1 -H, xedges, yedges = numpy.histogram2d(x=closest[0], y=closest[1], bins=41, range=[[-w, +w], [-h, +h]]) -fig = px.imshow(H, x=(xedges[: -1]+xedges[1: ])*0.5, y=(yedges[: -1]+yedges[1: ])*0.5) -#fig.show() -fig.write_image(generated_data + "/evaluation1_1.png") - -r = 6 -idx = 0 -shapes = [dict(x=row[0], y=row[1], sigma=r, color='green') - for row in true_data[idx][:, [3, 4]]] -shapes += [dict(x=spot[0], y=spot[1], sigma=r, color='red') - for spot in spots[idx]] -scopyon.Image(numpy.load(inputpath / "images{:03d}.npy".format(num_samples - 1))[idx]).show(shapes=shapes) +# import plotly.express as px +# w = h = 1 +# H, xedges, yedges = numpy.histogram2d(x=closest[0], y=closest[1], bins=41, range=[[-w, +w], [-h, +h]]) +# fig = px.imshow(H, x=(xedges[: -1]+xedges[1: ])*0.5, y=(yedges[: -1]+yedges[1: ])*0.5) +# fig.show() +# fig.write_image(str(artifacts / "heatmap1.png")) + +# r = 6 +# idx = 0 +# shapes = [dict(x=row[0], y=row[1], sigma=r, color='green') +# for row in true_data[idx][:, [3, 4]]] +# shapes += [dict(x=spot[0], y=spot[1], sigma=r, color='red') +# for spot in spots[idx]] +# scopyon.Image(numpy.load(generation_artifacts / "images{:03d}.npy".format(num_samples - 1))[idx]).show(shapes=shapes) r = rates[: 2].sum() / rates[2: ].sum() miss_count = rates[1] / rates[: 2].sum() @@ -133,6 +135,10 @@ log_metric("miss_count", miss_count) log_metric("missing", missing) -#log_artifacts("./artifacts") -log_artifacts(generated_data) +#XXX: THERE + +log_artifacts(str(artifacts)) mlflow.end_run() + +import shutil +shutil.rmtree(str(artifacts)) diff --git a/generation.py b/generation.py index c3d8244..f066cfc 100755 --- a/generation.py +++ b/generation.py @@ -1,34 +1,27 @@ -# -*- coding: utf-8 -*- -"""generation.ipynb - -Automatically generated by Colaboratory. - -Original file is located at - https://colab.research.google.com/github/ecell/bioimage_workflows/blob/master/generation.ipynb - -!pip uninstall -y scopyon -!pip install git+https://github.com/ecell/scopyon -!pip freeze | grep scopyon -""" - import argparse +import pathlib -"""Prepare for generating inputs.""" +import mlflow +from mlflow import log_metric, log_param, log_artifacts + +entrypoint = "generation" parser = argparse.ArgumentParser(description='generation step') +parser.add_argument('--seed', type=int, default=123) +parser.add_argument('--interval', type=float, default=33e-3) parser.add_argument('--num_samples', type=int, default=1) parser.add_argument('--num_frames', type=int, default=5) parser.add_argument('--exposure_time', type=float, default=0.033) args = parser.parse_args() -import mlflow -foo = mlflow.start_run(run_name="generation") +active_run = mlflow.start_run() +mlflow.set_tag("mlflow.runName", entrypoint) +seed = args.seed +interval = args.interval num_samples = args.num_samples num_frames = args.num_frames exposure_time = args.exposure_time -seed = 123 -interval = 33.0e-3 Nm = [100, 100, 100] Dm = [0.222e-12, 0.032e-12, 0.008e-12] transmat = [ @@ -36,16 +29,14 @@ [0.5, 0.0, 0.2], [0.0, 1.0, 0.0]] -from mlflow import log_metric, log_param, log_artifacts -log_param("seed", seed) -log_param("num_samples", num_samples) -log_param("num_frames", num_frames) -log_param("exposure_time", exposure_time) -log_param("interval", interval) +for key, value in vars(args).items(): + log_param(key, value) -#nproc = 8 +import tempfile +artifacts = pathlib.Path(tempfile.mkdtemp()) / "artifacts" +artifacts.mkdir(parents=True, exist_ok=True) -# !pip install mlflow +#XXX: HERE import numpy rng = numpy.random.RandomState(seed) @@ -59,25 +50,16 @@ L_2 = config.default.detector.image_size[0] * pixel_length * 0.5 L_2 -#config.environ.processes = nproc - timepoints = numpy.linspace(0, interval * num_frames, num_frames + 1) ndim = 2 -import pathlib -runid = foo.info.run_id -artifactsPath = "/tmp/" + str(runid) + "/artifacts" -artifacts = pathlib.Path(artifactsPath) -artifacts.mkdir(parents=True, exist_ok=True) -log_param("artifactsPath", artifactsPath) - config.save(artifacts / 'config.yaml') for i in range(num_samples): samples = scopyon.sample(timepoints, N=Nm, lower=-L_2, upper=+L_2, ndim=ndim, D=Dm, transmat=transmat, rng=rng) inputs = [(t, numpy.hstack((points[:, : ndim], points[:, [ndim + 1]], numpy.ones((points.shape[0], 1), dtype=numpy.float64)))) for t, points in zip(timepoints, samples)] ret = list(scopyon.generate_images(inputs, num_frames=num_frames, config=config, rng=rng, full_output=True)) - + inputs_ = [] for t, data in inputs: inputs_.extend(([t] + list(row) for row in data)) @@ -85,6 +67,7 @@ numpy.save(artifacts / f"inputs{i:03d}.npy", inputs_) numpy.save(artifacts / f"images{i:03d}.npy", numpy.array([img.as_array() for img, infodict in ret])) + ret[0][0].save(artifacts / f"image{i:03d}_000.png") true_data = [] for t, (_, infodict) in zip(timepoints, ret): @@ -92,7 +75,10 @@ true_data = numpy.array(true_data) numpy.save(artifacts / f"true_data{i:03d}.npy", true_data) -#!ls ./artifacts +#XXX: THERE -log_artifacts(artifactsPath) +log_artifacts(str(artifacts)) mlflow.end_run() + +import shutil +shutil.rmtree(str(artifacts)) diff --git a/main.py b/main.py index 41f2d43..48ba419 100644 --- a/main.py +++ b/main.py @@ -1,76 +1,13 @@ -import subprocess -import mlflow -from mlflow.utils import mlflow_tags -from mlflow.entities import RunStatus -from mlflow.utils.logging_utils import eprint - -from mlflow.tracking.fluent import _get_experiment_id -from mlflow import log_metric, log_param, log_artifacts import pathlib import argparse +import itertools -# _already_ran and _get_or_run code from bellow. -# mlflow/main.py at master · mlflow/mlflow -# https://github.com/mlflow/mlflow/blob/master/examples/multistep_workflow/main.py -# -# modify little bit at compare param with force string -def _already_ran(entry_point_name, parameters, git_commit, experiment_id=None): - """Best-effort detection of if a run with the given entrypoint name, - parameters, and experiment id already ran. The run must have completed - successfully and have at least the parameters provided. - """ - experiment_id = experiment_id if experiment_id is not None else _get_experiment_id() - client = mlflow.tracking.MlflowClient() - all_run_infos = reversed(client.list_run_infos(experiment_id)) - for run_info in all_run_infos: - full_run = client.get_run(run_info.run_id) - tags = full_run.data.tags - if tags.get(mlflow_tags.MLFLOW_PROJECT_ENTRY_POINT, None) != entry_point_name: - continue - match_failed = False - for param_key, param_value in parameters.items(): - run_value = full_run.data.params.get(param_key) - if str(run_value) != str(param_value): - match_failed = True - break - if match_failed: - continue - - if run_info.to_proto().status != RunStatus.FINISHED: - eprint( - ("Run matched, but is not FINISHED, so skipping " "(run_id=%s, status=%s)") - % (run_info.run_id, run_info.status) - ) - continue - - previous_version = tags.get(mlflow_tags.MLFLOW_GIT_COMMIT, None) - if git_commit != previous_version: - eprint( - ( - "Run matched, but has a different source version, so skipping " - "(found=%s, expected=%s)" - ) - % (previous_version, git_commit) - ) - continue - return client.get_run(run_info.run_id) - eprint("No matching run has been found.") - return None - - -# TODO(aaron): This is not great because it doesn't account for: -# - changes in code -# - changes in dependant steps -def _get_or_run(entrypoint, parameters, git_commit, use_cache=True): - existing_run = _already_ran(entrypoint, parameters, git_commit) - if use_cache and existing_run: - print("Found existing run for entrypoint=%s and parameters=%s" % (entrypoint, parameters)) - return existing_run - print("Launching new run for entrypoint=%s and parameters=%s" % (entrypoint, parameters)) - submitted_run = mlflow.run(".", entrypoint, parameters=parameters) - return mlflow.tracking.MlflowClient().get_run(submitted_run.run_id) +import mlflow +# from mlflow.utils import mlflow_tags +from mlflow import log_metric, log_param, log_artifacts +from mlflow_utils import _get_or_run -"""Prepare for generating inputs.""" +entrypoint = "main" parser = argparse.ArgumentParser(description='analysis1 step') parser.add_argument('--threshold', type=float, default=50.0) parser.add_argument('--min_sigma', type=int, default=1) @@ -83,37 +20,17 @@ def _get_or_run(entrypoint, parameters, git_commit, use_cache=True): num_samples = args.num_samples num_frames = args.num_frames -with mlflow.start_run(run_name="main", nested=True) as active_run: - # log param - log_param("threshold", threshold) - log_param("min_sigma", min_sigma) - log_param("num_samples", num_samples) - log_param("num_frames", num_frames) - # artifacts - #artifacts = pathlib.Path("./artifacts") - #artifacts.mkdir(parents=True, exist_ok=True) - # check git version - git_commit = active_run.data.tags.get(mlflow_tags.MLFLOW_GIT_COMMIT) - # generation - generation_run = _get_or_run("generation", {"num_samples":num_samples, "num_frames":num_frames}, git_commit) - artifactsPath = generation_run.data.params["artifactsPath"] - #generation_run = mlflow.run(".", "generation", parameters={"num_samples":num_samples, "num_frames":num_frames}) - # analysis1 - analysis1_run = _get_or_run("analysis1", {"generated_data":artifactsPath, "threshold":threshold, "min_sigma":min_sigma, "num_samples":num_samples, "num_frames":num_frames}, git_commit) - #analysis1_run = mlflow.run(".", "analysis1", parameters={"threshold":threshold, "num_samples":num_samples}) - # analysis2 - analysis2_run = _get_or_run("analysis2", {"threshold":threshold, "num_samples":num_samples, "num_frames":num_frames}, git_commit) - #analysis2_run = mlflow.run(".", "analysis2", parameters={"threshold":threshold, "num_samples":num_samples}) +with mlflow.start_run(nested=True) as active_run: + git_commit = active_run.data.tags.get("mlflow.source.git.branch") + mlflow.set_tag("mlflow.runName", entrypoint) + for key, value in vars(args).items(): + log_param(key, value) + + generation_run = _get_or_run("generation", {"num_samples": num_samples, "num_frames": num_frames}, git_commit) + analysis1_run = _get_or_run("analysis1", {"generation": generation_run.info.run_id, "threshold": threshold, "min_sigma": min_sigma}, git_commit) + analysis2_run = _get_or_run("analysis2", {"generation": generation_run.info.run_id, "analysis1": analysis1_run.info.run_id}, git_commit) + evaluation1_run = _get_or_run("evaluation1", {"generation": generation_run.info.run_id, "analysis1": analysis1_run.info.run_id}, git_commit) -# #log_artifacts("./artifacts") - # evaluation1 - evaluation1_run = _get_or_run("evaluation1", {"threshold":threshold, "num_samples":num_samples, "num_frames":num_frames}, git_commit) - #evaluation1_run = mlflow.run(".", "evaluation1", parameters={"threshold":threshold, "num_samples":num_samples}) - - log_metric("x_mean", float(evaluation1_run.data.metrics["x_mean"])) - log_metric("y_mean", float(evaluation1_run.data.metrics["y_mean"])) - log_metric("x_std", float(evaluation1_run.data.metrics["x_std"])) - log_metric("y_std", float(evaluation1_run.data.metrics["y_std"])) - log_metric("r", float(evaluation1_run.data.metrics["r"])) - log_metric("miss_count", float(evaluation1_run.data.metrics["miss_count"])) - log_metric("missing", float(evaluation1_run.data.metrics["missing"])) + for run_obj in (generation_run, analysis1_run, analysis2_run, evaluation1_run): + for key, value in run_obj.data.metrics.items(): + log_metric(key, value) diff --git a/mlflow_utils.py b/mlflow_utils.py new file mode 100644 index 0000000..89352c4 --- /dev/null +++ b/mlflow_utils.py @@ -0,0 +1,66 @@ +import mlflow + +from mlflow.tracking.fluent import _get_experiment_id +from mlflow.utils import mlflow_tags +from mlflow.entities import RunStatus +from mlflow.utils.logging_utils import eprint + +# _already_ran and _get_or_run code from bellow. +# mlflow/main.py at master · mlflow/mlflow +# https://github.com/mlflow/mlflow/blob/master/examples/multistep_workflow/main.py +# +# modify little bit at compare param with force string +def _already_ran(entry_point_name, parameters, git_commit, experiment_id=None): + """Best-effort detection of if a run with the given entrypoint name, + parameters, and experiment id already ran. The run must have completed + successfully and have at least the parameters provided. + """ + experiment_id = experiment_id if experiment_id is not None else _get_experiment_id() + client = mlflow.tracking.MlflowClient() + all_run_infos = reversed(client.list_run_infos(experiment_id)) + for run_info in all_run_infos: + full_run = client.get_run(run_info.run_id) + tags = full_run.data.tags + if tags.get(mlflow_tags.MLFLOW_PROJECT_ENTRY_POINT, None) != entry_point_name: + continue + match_failed = False + for param_key, param_value in parameters.items(): + run_value = full_run.data.params.get(param_key) + if str(run_value) != str(param_value): + match_failed = True + break + if match_failed: + continue + + if run_info.to_proto().status != RunStatus.FINISHED: + eprint( + ("Run matched, but is not FINISHED, so skipping " "(run_id=%s, status=%s)") + % (run_info.run_id, run_info.status) + ) + continue + + previous_version = tags.get(mlflow_tags.MLFLOW_GIT_COMMIT, None) + if git_commit != previous_version: + eprint( + ( + "Run matched, but has a different source version, so skipping " + "(found=%s, expected=%s)" + ) + % (previous_version, git_commit) + ) + continue + return client.get_run(run_info.run_id) + eprint("No matching run has been found.") + return None + +# TODO(aaron): This is not great because it doesn't account for: +# - changes in code +# - changes in dependant steps +def _get_or_run(entrypoint, parameters, git_commit, use_cache=True): + existing_run = _already_ran(entrypoint, parameters, git_commit) + if use_cache and existing_run: + print("Found existing run for entrypoint=%s and parameters=%s" % (entrypoint, parameters)) + return existing_run + print("Launching new run for entrypoint=%s and parameters=%s" % (entrypoint, parameters)) + submitted_run = mlflow.run(".", entrypoint, parameters=parameters) + return mlflow.tracking.MlflowClient().get_run(submitted_run.run_id)