关于Dask与PostgreSQL并行分布式读取数据的技术问询
问题
- 我尝试用Dask的
dd.read_sql_table()从AWS PostgreSQL数据库读取数据,传入完整的SQLAlchemy URI作为连接参数,想问这个方法是并行从数据库拉取数据,还是先拉取为Pandas DataFrame再由Dask分配到npartitions? - 我了解dask-sql,但
Context()需要先传入DataFrame作为输入,不管是否懒加载,都得先把数据加载到脚本里才能传入,这不符合我的需求。有没有既能通过多线程加速表查询,又能将数据分配到npartitions的方法?或者Dask在底层已经实现了查询的并行化? - 我自己构思了一段代码实现思路,想知道是否可行:
def get_dd_from_db(df_chunk): dd_from_pool = dd.from_pandas(df=df_chunk) return dd_from_pool def read_sql_chunks(table_name, con, chunksize): with ThreadPoolExecutor() as executor: futures = [] for chunk in pd.read_sql_table(table_name= table_name, con= con, chunksize= chunksize): future = executor.submit(chunk) futures.append(future) return futures def get_pd_from_db(table_name, con, chunksize): df_iterator = read_sql_chunks(table_name= table_name, con= con, chunksize= chunksize) with ProcessPoolExecutor() as executor: list(executor.map(lambda chunk: get_dd_from_db(chunk), df_iterator)) return df
思路是通过一个函数分块并行读取数据为Pandas DataFrame,再由另一个函数分块并行写入Dask DataFrame的各个分区。
回答
关于dd.read_sql_table()的并行机制
dd.read_sql_table()是直接并行从数据库拉取数据,并非先全量拉取为Pandas DataFrame再拆分分区。它的核心逻辑是:
- 根据你指定的
index_col(分区键,建议用有序的主键或整数列)或npartitions参数,自动将原表拆分为多个范围查询; - 每个分区对应一个独立的SQL查询任务,由Dask的任务调度器并行执行这些查询;
- 最终直接生成带有指定分区数的Dask DataFrame,全程无需将全量数据加载到内存。
关于多线程加速与分区的解决方案
Dask底层已经实现了SQL查询的并行化,完全不需要手动用线程池/进程池实现复杂的分块逻辑:
- 直接使用
dd.read_sql_table()的正确姿势
只需指定合适的分区参数即可实现并行读取,示例代码:
import dask.dataframe as dd from sqlalchemy import create_engine # 构建SQLAlchemy连接 con = f'{dialect}+{driver}://{username}:{password}@{host}:{port}/{database}' engine = create_engine(con) # 推荐:用有序列作为分区键(如自增主键),提升并行效率 ddf = dd.read_sql_table( table_name="your_target_table", con=engine, index_col="id", # 替换为你的表中适合分区的列 npartitions=8 # 根据数据量和集群资源调整分区数 ) # 若没有合适的分区键,也可仅指定npartitions(Dask会自动基于主键拆分) ddf = dd.read_sql_table( table_name="your_target_table", con=engine, npartitions=8 )
- dask-sql的正确用法(无需提前加载数据)
你对dask-sql的理解有误,它支持直接从数据库导入表,无需先将数据加载到内存:
from dask_sql import Context ctx = Context() # 直接从数据库导入表,底层复用Dask的并行读取能力 ctx.create_table( table_name="dask_table", con=con, # 传入SQLAlchemy URI或engine对象 schema="public", # 可选,指定数据库schema name="your_target_table" # 数据库中的源表名 ) # 后续可直接用SQL查询,返回的是懒加载的Dask DataFrame result_ddf = ctx.sql("SELECT * FROM dask_table WHERE status = 'active'")
对你构思代码的问题分析
你写的代码存在几个关键问题,无法实现预期的并行读取效果:
read_sql_chunks中executor.submit(chunk)是错误用法:submit需要传入可执行的函数,而非直接传chunk对象,会直接抛出异常;pd.read_sql_table的chunksize迭代是串行生成的,外层加线程池也无法实现并行读取,因为迭代器必须逐个生成chunk;get_dd_from_db将单个Pandas chunk转为Dask DataFrame,再合并的方式完全冗余,反而会增加不必要的开销,远不如Dask原生方法高效。
内容的提问来源于stack exchange,提问作者fCremer
相关产品推荐
相关产品推荐

