在Django中配置Prefect 2+Dask时任务反序列化失败的解决问询
Django + Prefect 2 + Dask 序列化报错解决办法
问题描述
在Django环境中配置Prefect 2与Dask运行任务时,简单示例可正常执行,但复杂示例运行失败,报错如下:
distributed.worker - ERROR - Could not deserialize task ****-****-****-**** Traceback (most recent call last): File "/usr/local/lib/python3.8/site-packages/distributed/worker.py", line 2185, in execute function, args, kwargs = await self._maybe_deserialize_task(ts) File "/usr/local/lib/python3.8/site-packages/distributed/worker.py", line 2158, in _maybe_deserialize_task function, args, kwargs = _deserialize(*ts.run_spec) File "/usr/local/lib/python3.8/site-packages/distributed/worker.py", line 2835, in _deserialize kwargs = pickle.loads(kwargs) File "/usr/local/lib/python3.8/site-packages/distributed/protocol/pickle.py", line 73, in loads return pickle.loads(x) File "/tmp/tmpb11md_6wprefect/my_project/models/__init__.py", line 2, in <module> from my_project.models.my_model import * File "/tmp/tmpb11md_6wprefect/my_project/models/api_key.py", line 14, in <module> class ApiKey(models.Model): File "/usr/local/lib/python3.8/site-packages/django/db/models/base.py", line 108, in __new__ app_config = apps.get_containing_app_config(module) File "/usr/local/lib/python3.8/site-packages/django/apps/registry.py", line 253, in get_containing_app_config self.check_apps_ready() File "/usr/local/lib/python3.8/site-packages/django/apps/registry.py", line 136, in check_apps_ready raise AppRegistryNotReady("Apps aren't loaded yet.") django.core.exceptions.AppRegistryNotReady: Apps aren't loaded yet.
项目未采用Django多App结构,而是通过model目录下的__init__.py统一导入模型。此前在其他场景中,会在任务开头初始化Django环境规避错误,代码如下:
import datetime import xyz import django import os os.environ.setdefault( "DJANGO_SETTINGS_MODULE", os.environ.get("DJANGO_SETTINGS_MODULE", "core.settings"), ) django.setup() # 现在导入模型 from my_project.models import MyModel
但该方法在Dask反序列化时失效,Dask会先导入model目录,再执行任务内的初始化代码。
解决办法
1. 延迟模型导入至任务执行阶段
不在任务函数顶层导入模型,将导入语句放在django.setup()之后的任务执行代码中,确保Django环境初始化完成后再加载模型:
import os import django from prefect import task @task def my_dask_task(): # 初始化Django环境 os.environ.setdefault( "DJANGO_SETTINGS_MODULE", os.environ.get("DJANGO_SETTINGS_MODULE", "core.settings"), ) django.setup() # 在这里导入模型 from my_project.models import MyModel # 执行任务逻辑 queryset = MyModel.objects.all() # ...其他操作
2. 配置Dask Worker启动时初始化Django环境
让每个Dask Worker启动时完成Django环境初始化,避免任务反序列化时触发模型导入。
创建dask_django_init.py初始化脚本:
import os import django os.environ.setdefault( "DJANGO_SETTINGS_MODULE", os.environ.get("DJANGO_SETTINGS_MODULE", "core.settings"), ) django.setup()
启动Dask Worker时指定预加载脚本:
dask-worker scheduler-address:8786 --preload dask_django_init.py
或在Prefect的Dask集群配置中指定:
from prefect_dask import DaskTaskRunner task_runner = DaskTaskRunner( cluster_kwargs={ "worker_kwargs": { "preload": ["dask_django_init.py"] } } ) @flow(task_runner=task_runner) def my_flow(): # ...流程逻辑
3. 修改模型导入方式,避免顶层触发Django加载
调整models/__init__.py的导入逻辑,改用延迟导入方式:
# 不要直接导入模型,定义函数在需要时获取 def get_api_key_model(): from my_project.models.api_key import ApiKey return ApiKey
在任务中通过函数获取模型:
@task def my_task(): django.setup() ApiKey = get_api_key_model() # ...操作模型
4. 使用Prefect Django集成简化初始化
安装Prefect Django集成工具:
pip install prefect-django
使用setup_django装饰器自动完成环境初始化:
from prefect import task from prefect_django import setup_django @task @setup_django def my_task(): from my_project.models import MyModel # ...任务逻辑
该装饰器会确保任务执行前完成Django环境初始化,同时规避反序列化时的导入问题。
内容的提问来源于stack exchange,提问作者Tim
相关产品推荐
相关产品推荐

