Slackbot
08/14/2023, 3:35 PMJian Shen Yap
08/15/2023, 7:18 PMMichael Wang
08/15/2023, 8:08 PMdagster-ml
- src
- my_dag
- assets
- __init__.py: loads assets from training_assets.py, preprocessing_assets.py, upload_assets.py
- preprocessing_assets.py: loads dbt assets into dagster
- training_assets.py: generates sklearn pipeline model
- upload_assets.py: uploads model into bentoml model registry
- __init__.py: creates definitions to be loaded into Dagster
Here is the error when I try to invoke the model locally:
local_runner = bentoml_model.to_runner()
local_runner.init_local()
print(local_runner.predict.run(pd.DataFrame(SAMPLE_PAYLOAD)))Michael Wang
08/15/2023, 8:10 PMjoblib.externals.loky.process_executor.BrokenProcessPool: A task has failed to un-serialize. Please ensure that the arguments of the function are all picklable.
Stack Trace:
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/dagster/_core/execution/plan/utils.py", line 54, in op_execution_error_boundary
yield
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/dagster/_utils/__init__.py", line 445, in iterate_with_context
next_output = next(iterator)
^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/dagster/_core/execution/plan/compute_generator.py", line 124, in _coerce_op_compute_fn_to_iterator
result = invoke_compute_fn(
^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/dagster/_core/execution/plan/compute_generator.py", line 118, in invoke_compute_fn
return fn(context, **args_to_pass) if context_arg_provided else fn(**args_to_pass)
^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/Documents/dagster-ml/src/my_dag/assets/upload_assets.py", line 105, in upload_to_model_registry
print(local_runner.predict.run(pd.DataFrame(SAMPLE_PAYLOAD)))
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/bentoml/_internal/runner/runner.py", line 52, in run
return self.runner._runner_handle.run_method(self, *args, **kwargs)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/bentoml/_internal/runner/runner_handle/local.py", line 48, in run_method
return getattr(self._runnable, __bentoml_method.name)(*args, **kwargs)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/bentoml/_internal/runner/runnable.py", line 140, in method
return self.func(obj, *args, **kwargs)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/bentoml/_internal/frameworks/sklearn.py", line 190, in _run
return getattr(self.model, method_name)(input_data)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/sklearn/pipeline.py", line 507, in predict
Xt = transform.transform(Xt)
^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/sklearn/utils/_set_output.py", line 140, in wrapped
data_to_wrap = f(self, X, *args, **kwargs)
^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/sklearn/compose/_column_transformer.py", line 816, in transform
Xs = self._fit_transform(
^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/sklearn/compose/_column_transformer.py", line 670, in _fit_transform
return Parallel(n_jobs=self.n_jobs)(
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/sklearn/utils/parallel.py", line 65, in __call__
return super().__call__(iterable_with_config)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/joblib/parallel.py", line 1952, in __call__
return output if self.return_generator else list(output)
^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/joblib/parallel.py", line 1595, in _get_outputs
yield from self._retrieve()
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/joblib/parallel.py", line 1699, in _retrieve
self._raise_error_fast()
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/joblib/parallel.py", line 1734, in _raise_error_fast
error_job.get_result(self.timeout)
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/joblib/parallel.py", line 736, in get_result
return self._return_or_raise()
^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/joblib/parallel.py", line 754, in _return_or_raise
raise self._result
The above exception was caused by the following exception:
joblib.externals.loky.process_executor._RemoteTraceback:
"""
Traceback (most recent call last):
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/joblib/externals/loky/process_executor.py", line 426, in _process_worker
call_item = call_queue.get(block=True, timeout=timeout)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/multiprocessing/queues.py", line 122, in get
return _ForkingPickler.loads(res)
^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/Documents/dagster-ml/src/my_dag/__init__.py", line 10, in <module>
from src.my_dag.assets import (
File "/Users/michael.wang/Documents/dagster-ml/src/my_dag/assets/__init__.py", line 3, in <module>
from src.my_dag.assets.preprocessing_assets import dbt_assets
File "/Users/michael.wang/Documents/dagster-ml/src/my_dag/assets/preprocessing_assets.py", line 24, in <module>
dbt_assets = load_assets_from_dbt_project(
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/dagster_dbt/asset_defs.py", line 606, in load_assets_from_dbt_project
manifest, cli_output = _load_manifest_for_project(
^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/dagster_dbt/asset_defs.py", line 84, in _load_manifest_for_project
cli_output = execute_cli(
^^^^^^^^^^^^
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/dagster_dbt/core/utils.py", line 235, in execute_cli
for event in _core_execute_cli(
File "/Users/michael.wang/miniconda3/envs/dagster/lib/python3.11/site-packages/dagster_dbt/core/utils.py", line 160, in _core_execute_cli
raise DagsterDbtCliHandledRuntimeError(messages=messages)
dagster_dbt.errors.DagsterDbtCliHandledRuntimeError: Handled error in the dbt CLI (return code 1)
"""Michael Wang
08/15/2023, 8:12 PMdagster-ml/src/my_dag/__init__.py , which runs dagster-ml/src/my_dag/assets/preprocessing_assets.py , even though neither __init__.py nor preprocessing_assets.py are part of the model dependenciesJian Shen Yap
08/15/2023, 8:31 PMtraining_assets.py is where you train your sklearn model and save it as a bentoml model?Michael Wang
08/15/2023, 8:32 PMJian Shen Yap
08/15/2023, 8:35 PMMichael Wang
08/15/2023, 8:37 PM@asset()
def upload_to_model_registry(trained_model_results: sklearn.Pipeline):
model_name = 'my-model'
bentoml_model = bentoml.sklearn.save_model(
name=model_name,
model=trained_model_results
)Jian Shen Yap
08/15/2023, 9:10 PMbentoml models get <your model tag> -o yaml
at the same time you can visit the path where your model is saved, to see what is being saved along
bentoml models get <your model tag> -o pathMichael Wang
08/15/2023, 9:15 PMname: model-name
version: ihqj72byrsivxyrs
module: bentoml.sklearn
labels: {}
options: {}
metadata: {}
context:
framework_name: sklearn
framework_versions:
scikit-learn: 1.3.0
bentoml_version: 1.1.1
python_version: 3.11.4
signatures:
predict:
batchable: false
api_version: v1
creation_time: '2023-08-11T21:15:54.688807+00:00'Jian Shen Yap
08/15/2023, 9:19 PMload_model triggers the execution of your __init__.py
are you doing something like from my_dag.assets import xxxMichael Wang
08/15/2023, 9:21 PM__init__.py , which is the standard dagster design patternJian Shen Yap
08/15/2023, 9:30 PMload_model? i don't see it in your directory structure,Michael Wang
08/15/2023, 9:31 PMload_model but I do to_runner immediately after I call `save_model`:
@asset()
def upload_to_model_registry(trained_model_results: sklearn.Pipeline):
model_name = 'my-model'
bentoml_model = bentoml.sklearn.save_model(
name=model_name,
model=trained_model_results
)
local_runner = bentoml_model.to_runner()
local_runner.init_local()
print(local_runner.predict.run(pd.DataFrame(SAMPLE_PAYLOAD))) >>> ERRORJian Shen Yap
08/15/2023, 10:09 PMit looks like it’s packaging the entire directoryI can confirm this is not true. From the stacktrace, it looks like its some serialization issue when running a BentoML runner in dagster context.
Michael Wang
08/15/2023, 10:13 PMJian Shen Yap
08/15/2023, 10:15 PMloky backend on the runner sidesauyon
08/15/2023, 10:17 PMAaron Pham
08/15/2023, 10:20 PMAaron Pham
08/15/2023, 10:21 PMAaron Pham
08/15/2023, 10:23 PMto_runner directly after save model. Is there any reason why you are doing so?
https://github.com/bentoml/BentoML/blob/ced530b816235aefa2515141087fc3320896f423/src/bentoml/_internal/frameworks/sklearn.py#L189Aaron Pham
08/15/2023, 10:25 PMMichael Wang
08/17/2023, 11:42 PM@assets decorator adds that functions to the dag that dagster runs. It controls how data flows in and out of the graph and allows upstream assets to feed into downstream assets.
• I am calling to_runner directly after saving the model to replicate the code I saw on some of the bentoml docs. In production, I’d like to unit test the model runner to make sure input/output are what we expect before pushing the model out to Yatai.
I don’t think it’s the @asset decorator that is causing the serialization issue. I created a .py file without using @asset and still got same error from before:
bentoml_model = bentoml.sklearn.save_model(
name='my-model',
model=model
)
runner = bentoml_model.to_runner()
runner.init_local()
runner.predict.run(pd.DataFrame([payload]))