Настройка
Код: Выделить всё
import os
os.environ["MODIN_ENGINE"] = "dask"
from dask.distributed import Client, LocalCluster
cluster = LocalCluster(
n_workers=1,
threads_per_worker=8,
processes=False, # same-process threads to avoid cross-process pickle
memory_limit="8GB",
scheduler_port=0,
dashboard_address=None,
)
client = Client(cluster)
import modin.pandas as mpd
df = mpd.DataFrame({"a": range(100_000)})
result = df.sort_values("a") # triggers the error
Код: Выделить всё
INFO distributed.protocol.pickle:pickle.py:97 Failed to deserialize
Traceback (most recent call last):
File ".../distributed/protocol/pickle.py", line 95, in loads
return pickle.loads(x)
File ".../twisted/persisted/styles.py", line 68, in unpickleMethod
methodFunction = _methodFunction(im_class, im_name)
File ".../twisted/persisted/styles.py", line 48, in _methodFunction
methodObject = getattr(classObject, methodName)
AttributeError: type object 'ABCMeta' has no attribute 'deploy_axis_func'
ERROR distributed.scheduler:scheduler.py:4973 Error during deserialization of the task graph.
This frequently occurs if the Scheduler and Client have different environments.
Мы запускаем это внутри приложения Django, где фоновые рабочие процессы выполняют конвейеры преобразования данных (сортировка, объединение, сведение) для больших наборов данных, которые могут состоять из десятков миллионов строк. Мы перешли с простых панд на Modin, чтобы распараллелить эти операции. Поскольку рабочие уже хранят DataFrames в памяти внутри процесса Django, мы выбралиprocesses=False, чтобы рабочие Dask использовали один и тот же контекст процесса без необходимости межпроцессной сериализации. Ошибка возникает, как только какая-либо операция Modin (сортировка, группировка и т. д.) отправляется в кластер.
Насколько я понимаю, Deploy_axis_func — это метод внутреннего класса DataFrame Modin. Иерархия классов Modin использует ABCMeta в качестве метакласса. Когда распределенный Dask сериализует связанную ссылку на метод, unpickleMethod Twisted пытается восстановить ее, вызывая getattr(ABCMeta, 'deploy_axis_func') — просматривая метакласс вместо фактического класса — и терпит неудачу.
Это происходит даже сprocesses=False, потому что dask.distributed по-прежнему сериализует граф задачи между клиентом и внутрипроцессный планировщик через протокол Twisted, независимо от того, являются ли рабочие потоками или отдельными процессами.
Версии
- Python 3.13
- modin 0.37.1
- dask / распределенный 2025.x
- twisted (последний)
Есть ли способ использовать бэкэнд Modin Dask с параллелизмом потоков, полностью обходя dask.distributed — например, с помощью собственного многопоточного планировщика Dask (
Код: Выделить всё
dask.config.set(scheduler='threads'))? Или есть известный обходной путь для несовместимости ABCMeta