使用Dask read_sql_query查询AWS Redshift触发ProgrammingError
使用Dask read_sql_query查询AWS Redshift时触发语法错误
尝试用Dask的read_sql_query方法查询AWS Redshift,运行代码时触发redshift_connector.ProgrammingError,提示SQL语句中SELECT附近存在语法错误。相关代码及报错回溯如下:
相关代码
import dask.dataframe as dd from config import * host=os.environ['host'] database=os.environ['database'] port='5439' user=os.environ['user'] password=os.environ['password'] conn_str = f'redshift+redshift_connector://{user}:{password}@{host}:{port}/{database}' # Query table using dask dataframe query = text(''' SELECT * FROM tbl WHERE type = :type AND created_at >= :start_date AND created_at <= :end_date ''') params = { 'type': 'type1', 'start_date': '2023-01-01 00:00:00', 'end_date': '2023-06-30 00:00:00' } selectable = select('*').select_from(query).params(**params) print(type(selectable)) df = dd.read_sql_query(selectable, conn_str, index_col = 'id')
报错回溯
<class 'sqlalchemy.sql.selectable.Select'> Unexpected exception formatting exception. Falling back to standard exception Traceback (most recent call last): File "/local_disk0/.ephemeral_nfs/cluster_libraries/python/lib/python3.9/site-packages/redshift_connector/core.py", line 1631, in execute ps = cache["ps"][key] KeyError: ('SELECT * \nFROM \n SELECT * \n FROM pmf\n WHERE event_type = %s\n AND created_at >= %s\n AND created_at <= %s \n \n LIMIT 5', ((<RedshiftOID.UNKNOWN: 705>, 0, <function text_out at 0x7f7a348d3790>), (<RedshiftOID.UNKNOWN: 705>, 0, <function text_out at 0x7f7a348d3790>), (<RedshiftOID.UNKNOWN: 705>, 0, <function text_out at 0x7f7a348d3790>))) During handling of the above exception, another exception occurred: Traceback (most recent call last): File "/local_disk0/.ephemeral_nfs/cluster_libraries/python/lib/python3.9/site-packages/sqlalchemy/engine/base.py", line 1900, in _execute_context self.dialect.do_execute( File "/local_disk0/.ephemeral_nfs/cluster_libraries/python/lib/python3.9/site-packages/sqlalchemy/engine/default.py", line 736, in do_execute cursor.execute(statement, parameters) File "/local_disk0/.ephemeral_nfs/cluster_libraries/python/lib/python3.9/site-packages/redshift_connector/cursor.py", line 240, in execute self._c.execute(self, operation, args) File "/local_disk0/.ephemeral_nfs/cluster_libraries/python/lib/python3.9/site-packages/redshift_connector/core.py", line 1701, in execute self.handle_messages(cursor) redshift_connector.error.ProgrammingError: {'S': 'ERROR', 'C': '42601', 'M': 'syntax error at or near "SELECT"', 'P': '30', 'F': '/home/ec2-user/padb/src/pg/src/backend/parser/parser_scan.l', 'L': '732', 'R': 'yyerror'} The above exception was the direct cause of the following exception: Traceback (most recent call last): File "/databricks/python/lib/python3.9/site-packages/IPython/core/interactiveshell.py", line 3378, in run_code exec(code_obj, self.user_global_ns, self.user_ns) File "<command-2539550446659032>", line 19, in <module> df = dd.read_sql_query(selectable, conn_str, index_col = 'id') File "/local_disk0/.ephemeral_nfs/cluster_libraries/python/lib/python3.9/site-packages/dask/dataframe/io/sql.py", line 120, in read_sql_query head = pd.read_sql(q, engine, **kwargs) File "/local_disk0/.ephemeral_nfs/cluster_libraries/python/lib/python3.9/site-packages/redshift_connector/core.py", line 1969, in handle_messages raise self.error sqlalchemy.exc.ProgrammingError: (redshift_connector.error.ProgrammingError) {'S': 'ERROR', 'C': '42601', 'M': 'syntax error at or near "SELECT"', 'P': '30', 'F': '/home/ec2-user/padb/src/pg/src/backend/parser/parser_scan.l', 'L': '732', 'R': 'yyerror'} [SQL: SELECT * FROM SELECT * FROM pmf WHERE event_type = %s AND created_at >= %s AND created_at <= %s LIMIT 5] [parameters: ('type1', '2023-01-01 00:00:00', '2023-06-30 00:00:00')] (Background on this error at: https://sqlalche.me/e/14/f405) During handling of the above exception, another exception occurred: Traceback (most recent call last): File "/databricks/python/lib/python3.9/site-packages/IPython/core/interactiveshell.py", line 1997, in showtraceback stb = self.InteractiveTB.structured_traceback( File "/databricks/python/lib/python3.9/site-packages/IPython/core/ultratb.py", line 1112, in structured_traceback return FormattedTB.structured_traceback( File "/databricks/python/lib/python3.9/site-packages/IPython/core/ultratb.py", line 1006, in structured_traceback return VerboseTB.structured_traceback( File "/databricks/python/lib/python3.9/site-packages/IPython/core/ultratb.py", line 859, in structured_traceback formatted_exception = self.format_exception_as_a_whole(etype, evalue, etb, number_of_lines_of_context, File "/databricks/python/lib/python3.9/site-packages/IPython/core/ultratb.py", line 812, in format_exception_as_a_whole frames.append(self.format_record(r)) File "/databricks/python/lib/python3.9/site-packages/IPython/core/ultratb.py", line 730, in format_record result += ''.join(_format_traceback_lines(frame_info.lines, Colors, self.has_colors, lvals)) File "/local_disk0/.ephemeral_nfs/cluster_libraries/python/lib/python3.9/site-packages/stack_data/utils.py", line 145, in cached_property_wrapper value = obj.__dict__[self.func.__name__] = self.func(obj) return only( File "/local_disk0/.ephemeral_nfs/cluster_libraries/python/lib/python3.9/site-packages/executing/executing.py", line 164, in only raise NotOneValueFound('Expected one value, found 0') executing.executing.NotOneValueFound: Expected one value, found 0 ProgrammingError: (redshift_connector.error.ProgrammingError) {'S': 'ERROR', 'C': '42601', 'M': 'syntax error at or near "SELECT"', 'P': '30', 'F': '/home/ec2-user/padb/src/pg/src/backend/parser/parser_scan.l', 'L': '732', 'R': 'yyerror'} [SQL: SELECT * FROM SELECT * FROM tbl WHERE event_type = %s AND created_at >= %s AND created_at
问题分析
从报错里生成的SQL语句可以明确问题:
SELECT * FROM SELECT * FROM tbl WHERE event_type = %s AND created_at >= %s AND created_at <= %s LIMIT 5
这段SQL将原查询作为子查询嵌套,但子查询未加括号,违反了Redshift的SQL语法规则,直接导致语法错误。
错误根源是代码中用select('*').select_from(query)对已经是完整查询的text对象再次嵌套,完全属于多余操作。
修复方案
方案一:直接传递text查询和参数
dd.read_sql_query支持直接接收SQLAlchemy的text对象作为查询语句,同时通过params参数传递绑定变量,无需额外包装:
import dask.dataframe as dd from sqlalchemy import text import os host=os.environ['host'] database=os.environ['database'] port='5439' user=os.environ['user'] password=os.environ['password'] conn_str = f'redshift+redshift_connector://{user}:{password}@{host}:{port}/{database}' # 定义完整查询语句 query = text(''' SELECT * FROM tbl WHERE type = :type AND created_at >= :start_date AND created_at <= :end_date ''') params = { 'type': 'type1', 'start_date': '2023-01-01 00:00:00', 'end_date': '2023-06-30 00:00:00' } # 直接传递query和params给read_sql_query df = dd.read_sql_query(query, conn_str, index_col='id', params=params)
方案二:使用SQLAlchemy查询构建器
如果想用SQLAlchemy的ORM风格构建查询,应该直接绑定表对象,而非嵌套text查询:
import dask.dataframe as dd from sqlalchemy import create_engine, select, Table, MetaData import os host=os.environ['host'] database=os.environ['database'] port='5439' user=os.environ['user'] password=os.environ['password'] conn_str = f'redshift+redshift_connector://{user}:{password}@{host}:{port}/{database}' engine = create_engine(conn_str) # 加载表元数据 metadata = MetaData() tbl = Table('tbl', metadata, autoload_with=engine) # 构建查询 selectable = select(tbl).where( tbl.c.type == 'type1', tbl.c.created_at >= '2023-01-01 00:00:00', tbl.c.created_at <= '2023-06-30 00:00:00' ) # 执行查询 df = dd.read_sql_query(selectable, conn_str, index_col='id')
关键提示
- 当已经用
text()定义完整SQL语句时,不要用select()再次包装,否则会生成无效的嵌套SQL。 dd.read_sql_query兼容SQLAlchemy的text对象和Select对象,按需选择即可。
内容的提问来源于stack exchange,提问作者kms
相关产品推荐
相关产品推荐

