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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 12:02:01