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
相关产品推荐
相关产品推荐

