This message was deleted.
# ask-for-help
s
This message was deleted.
j
Hi @Michael Wang, do you have some example that you can share, i.e your code and your repository structure of the Dagster project
m
Sure I can share the repository structure:
Copy code
dagster-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:
Copy code
local_runner = bentoml_model.to_runner()
local_runner.init_local()
print(local_runner.predict.run(pd.DataFrame(SAMPLE_PAYLOAD)))
Copy code
joblib.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)
"""
Note that the dagster dbt cli error is from the model running
dagster-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 dependencies
j
i presume
training_assets.py
is where you train your sklearn model and save it as a bentoml model?
m
training_assets.py is where I train the model, and it gets passed into upload_assets.py to be saved as a bentoml model
j
could you share a code snippets on how are you saving the model?
m
Sure, here’s a simplified version of the asset in upload_assets.py:
Copy code
@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
        )
j
I dont think this is an intended behaviour, could you share the output of
bentoml 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 path
m
yaml:
Copy code
name: 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'
j
i think its most likely that the way you are invoking
load_model
triggers the execution of your
__init__.py
are you doing something like
from my_dag.assets import xxx
m
yes, we have some of these imports in the
__init__.py
, which is the standard dagster design pattern
j
where are you trying to
load_model
? i don't see it in your directory structure,
m
I don’t call
load_model
but I do
to_runner
immediately after I call `save_model`:
Copy code
@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))) >>> ERROR
j
it looks like it’s packaging the entire directory
I can confirm this is not true. From the stacktrace, it looks like its some serialization issue when running a BentoML runner in dagster context.
m
Yep, agreed. There is no issue when pushing it out to Yatai and deploying it within a container, but something strange is happening when we load the model in the dagster context
j
@sauyon do you think you might be able to spot something here? error seems to be related to joblib's
loky
backend on the runner side
s
@Aaron Pham or @larme might know more off the top of their heads on frameworks, if not I'll get back to this tomorrow.
a
Seems like this @assets decorator is making this function not serialisable
I don’t know dagster throughout, but can you explain briefly what is this decorator is doing?
additionally, I don’t think you should call
to_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#L189
iirc loky is process-based parallelism, so it should be thread-safe
m
Sure, @Aaron Pham : • the
@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:
Copy code
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]))