You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Dask分布式Docker环境下SQLAlchemy查询报错:无法pickle 'weakref'对象

问题:Docker中Dask Worker执行SQLAlchemy查询时出现pickle弱引用错误

在Docker容器中运行Dask分布式应用,执行SQLAlchemy的read_sql_query语句时,worker抛出异常:Exception: cannot pickle 'weakref' object。该代码在本地worker中运行正常,仅在Docker环境下失败。

简化代码

from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy import Column, String, Integer, PrimaryKeyConstraint, select
from sqlalchemy.orm import aliased
import dask.dataframe as dd

Base = declarative_base()

class test_loans(Base):
    __tablename__ = 'test_loans'
    loan_id = Column(String)
    fico_score = Column(Integer)
    __table_args__ = (PrimaryKeyConstraint('fico_score'),)

t = aliased(test_loans)
stmt2 = select([t.loan_id, t.fico_score])
ddf = dd.read_sql_query(stmt2, con=db_str, index_col='fico_score', npartitions=20)
ddf.compute() # <----- 此处触发worker执行时失败

异常信息

2022-08-15 20:37:26,984 - distributed.protocol.pickle - INFO - Failed to serialize (<function apply at 0x7fcfee313160>, <function _read_sql_chunk at 0x7fcfbc7115e0>, [<sqlalchemy.sql.selectable.Select object at 0x7fcf6edc5430>, 'mariadb://user:xxxxx@host.docker.internal:3306/bank_0001', Empty DataFrame
Columns: [loan_id]
Index: []], (<class 'dict'>, [['engine_kwargs', (<class 'dict'>, [])], ['index_col', 'fico_score']])). Exception: cannot pickle 'weakref' object
2022-08-15 20:37:26,987 - distributed.protocol.core - CRITICAL - Failed to Serialize
Traceback (most recent call last):
  File "/opt/conda/lib/python3.8/site-packages/distributed/protocol/core.py", line 109, in dumps
    frames[0] = msgpack.dumps(msg, default=_encode_default, use_bin_type=True)
  File "/opt/conda/lib/python3.8/site-packages/msgpack/__init__.py", line 38, in packb
    return Packer(**kwargs).pack(o)
  File "msgpack/_packer.pyx", line 294, in msgpack._cmsgpack.Packer.pack
  File "msgpack/_packer.pyx", line 300, in msgpack._cmsgpack.Packer.pack
  File "msgpack/_packer.pyx", line 297, in msgpack._cmsgpack.Packer.pack
  File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack
  File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack
  File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack
  File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack
  File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack
  File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack
  File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack
  File "msgpack/_packer.pyx", line 285, in msgpack._cmsgpack.Packer._pack
  File "/opt/conda/lib/python3.8/site-packages/distributed/protocol/core.py", line 100, in _encode_default
    frames.extend(create_serialized_sub_frames(obj))
  File "/opt/conda/lib/python3.8/site-packages/distributed/protocol/core.py", line 60, in create_serialized_sub_frames
    sub_header, sub_frames = serialize_and_split(
  File "/opt/conda/lib/python3.8/site-packages/distributed/protocol/serialize.py", line 444, in serialize_and_split
    header, frames = serialize(x, serializers, on_error, context)
  File "/opt/conda/lib/python3.8/site-packages/distributed/protocol/serialize.py", line 266, in serialize
    return serialize(
  File "/opt/conda/lib/python3.8/site-packages/distributed/protocol/serialize.py", line 316, in serialize
    headers_frames = [
  File "/opt/conda/lib/python3.8/site-packages/distributed/protocol/serialize.py", line 317, in <listcomp>
    serialize(
  File "/opt/conda/lib/python3.8/site-packages/distributed/protocol/serialize.py", line 366, in serialize
    raise TypeError(msg, str(x)[:10000])
TypeError: ('Could not serialize object of type tuple', "('<function apply at 0x7fcfee313160>', '<function _read_sql_chunk at 0x7fcfbc7115e0>', [<sqlalchemy.sql.selectable.Select object at 0x7fcf6edc5430>, 'mariadb://xxxx:yyyy@host.docker.internal:3306/bank_0001', Empty DataFrame
Columns: [loan_id]
Index: []], (<class 'dict'>, [['engine_kwargs', (<class 'dict'>, [])], ['index_col', 'fico_score']]))")
2022-08-15 20:37:26,989 - distributed.comm.utils - INFO - Unserializable Message: [{'op': 'update-graph-hlg', 'hlg': {'layers': [{'__module__': 'dask.highlevelgraph', '__name__': 'MaterializedLayer', 'state': {'dsk': {\"('from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 2)\": <Serialize: ('read_sql_chunk-from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 2)>, \"('from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 8)\": <Serialize: ('read_sql_chunk-from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 8)>, \"('from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 5)\": <Serialize: ('read_sql_chunk-from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 5)>, \"('read_sql_chunk-from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 2)\": <Serialize: (subgraph_callable-8219b5cb-9a04-436c-ad29-d643f53e22d3, (<function apply at 0x7fcfee313160>, <function _read_sql_chunk at 0x7fcfbc7115e0>, [<sqlalchemy.sql.selectable.Select object at 0x7fcf6edc5430>, 'mariadb://xxxx:yyyy@host.docker.internal:3306/bank_0001', Empty DataFrame
Columns: [loan_id]
Index: []], (<class 'dict'>, [['engine_kwargs', (<class 'dict'>, [])], ['index_col', 'fico_score']]))>, \"('from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 1)\": <Serialize: ('read_sql_chunk-from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 1)>, \"('read_sql_chunk-from-delayed-45f566e88ca4c6d12a8197c6d219b27b', 5)\": <Serialize: (subgraph_callable-8219b5cb-9a04-436c-ad29-d643f53e22d3, (<function apply at 0x7fcfee313160>, <function _read_sql_chunk at 0x7fcfbc7115e0>, [<sqlalchemy.sql.selectable.Select object at 0x7fcf6edc5bb0>, 'mariadb://xxxx:yyyy@host.docker.internal:3306/bank_0001', Empty DataFrame

依赖包版本

dask                         2022.8.0
dask-glm                     0.2.0
dask-ml                      2022.5.27
dask-xgboost                 0.2.0
distributed                  2022.8.0
SQLAlchemy                   1.4.40
cloudpickle                  2.1.0

原因与解决方案

原因分析

这不是操作错误,核心问题是SQLAlchemy的Select对象内部包含无法被pickle序列化的弱引用对象:

  • 本地运行时,Dask任务无需跨进程/节点传输,直接在本地执行,不会触发序列化操作,因此正常运行;
  • Docker环境下,Dask需要将任务(包括Select对象)序列化后发送给远程Worker,此时弱引用对象无法被序列化,导致报错。

解决方案

1. 将Select对象转换为SQL字符串传递

这是最直接的解决方案,避免传递无法序列化的Select对象,改为传递原始SQL字符串:

stmt2 = select([t.loan_id, t.fico_score])
# 将Select语句编译为SQL字符串
sql_str = str(stmt2.compile(compile_kwargs={"literal_binds": True}))
# 使用SQL字符串执行查询
ddf = dd.read_sql_query(sql_str, con=db_str, index_col='fico_score', npartitions=20)

注意:literal_binds=True会直接将参数值绑定到SQL中,存在SQL注入风险。生产环境建议使用参数化查询,通过params传递参数:

sql_str = str(stmt2)
ddf = dd.read_sql_query(sql_str, con=db_str, index_col='fico_score', npartitions=20, params={...})

2. 升级依赖包

新版本的Dask、SQLAlchemy或cloudpickle可能修复了序列化兼容性问题,建议尝试升级到较新的稳定版本:

  • Dask/distributed:升级到2023.x及以上版本
  • SQLAlchemy:升级到2.x及以上版本
  • cloudpickle:升级到3.x及以上版本

3. 自定义序列化逻辑(不推荐)

通过扩展cloudpickle的序列化器,为SQLAlchemy的Select对象添加自定义序列化支持,但这种方式需要深入理解序列化机制,维护成本较高,仅作为备选方案。


内容的提问来源于stack exchange,提问作者ps0604

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 15:27:12