Dask read_sql_query计算时SQL列不存在错误排查求助
使用Dask的read_sql_query方法创建DataFrame后,常规数据操作可正常执行,但执行需要计算的操作(如透视表)时,触发SQL错误,提示找不到设置为索引的列。
相关代码
import pyodbc import pandas as pd from dask.dataframe import read_sql_query import dask as dk from sqlalchemy import sql import urllib from sqlalchemy import create_engine params = urllib.parse.quote_plus(conn_string) engine = create_engine("mssql+pyodbc:///?odbc_connect=%s" % params) url_str = f"mssql+pyodbc:///?odbc_connect=%s" % params def _remove_leading_select_from_query(query): if query.startswith("SELECT "): return query.replace("SELECT ", "", 1) else: return query def query_sql_to_dataframe(conn_string, query,index): try: conection = pyodbc.connect(conn_string) sa_query = sql.select(sql.text(_remove_leading_select_from_query(query))) df = read_sql_query(sa_query, url_str,index_col=index) except pyodbc.Error as ex: df = None finally: if conection : conection .close() return df # FT_DOC_ELETRONICO query = 'SELECT ROW_NUMBER() OVER(ORDER BY some_column ASC) AS row,* FROM dbo.my_table' idx = "row" doc_eletr = query_sql_to_dataframe(conn_string, query,index=idx)
错误信息
ProgrammingError: (pyodbc.ProgrammingError) ('42S22', "[42S22] [Microsoft][ODBC Driver 17 for SQL Server][SQL Server]Invalid column name 'row'. (207) (SQLExecDirectW); [42S22] [Microsoft][ODBC Driver 17 for SQL Server][SQL Server]Invalid column name 'row'. (207); [42S22] [Microsoft][ODBC Driver 17 for SQL Server][SQL Server]Statements could not be prepared (8180)")
[SQL: SELECT ROW_NUMBER() OVER(ORDER BY some_column ASC) AS row,* FROM dbo.my_table
WHERE row >= ? AND row <= ?]
[parameters: (1, 1181321)]
问题原因与解决方案
原因
Dask的read_sql_query为实现分区查询,会自动在原SQL后追加WHERE 索引列 >= ? AND 索引列 <= ?的过滤条件。但你的索引列row是SELECT子句中生成的计算列别名,SQL Server的执行顺序是先处理WHERE子句,再处理SELECT子句,因此WHERE无法直接引用SELECT中定义的别名,导致"找不到列"的错误。
解决方案
方案1:将原查询改为子查询,让索引列成为可引用的列
把生成row列的逻辑放到子查询中,外层查询直接引用该列,这样Dask追加的WHERE条件就能正确找到索引列:
# 修改后的查询语句 query = ''' SELECT * FROM ( SELECT ROW_NUMBER() OVER(ORDER BY some_column ASC) AS row,* FROM dbo.my_table ) AS subquery '''
同时简化query_sql_to_dataframe函数中的SQLAlchemy语句,避免多余的嵌套SELECT:
# 替换函数内的sa_query行 sa_query = sql.text(query)
方案2:使用表中已存在的原生列作为索引
如果目标表中有原生的主键、唯一键或连续值列(比如某个ID列),直接将该列指定为index_col,这样Dask生成的WHERE条件会直接引用表中存在的列,不会出现别名引用问题。
内容的提问来源于stack exchange,提问作者dsnishimura

