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

AWS Glue3 PySpark用psycopg2删除Aurora数据报col2不在列表错误

根本原因

这个错误和PostgreSQL语法、psycopg2执行逻辑没有任何关系,是PySpark侧的字段访问错误:
你传入的tuple_list不是你预期的原生Python元组列表,实际是PySpark Row对象构成的集合。报错栈明确显示错误发生在pyspark/sql/types.py的Row属性读取逻辑——你在execute_batch_delete方法第88行写了类似row.col2的属性访问代码,但当前处理的Row对象的schema里根本不存在名为col2的字段,所以直接抛出AttributeError,根本没走到SQL执行的步骤。

常见触发场景:

  • 生成待删除数据集的DataFrame/DynamicFrame时,select阶段漏选了对应字段
  • 字段名拼写错误:比如实际字段名是col_2、create_time、Col2(大小写不匹配),你硬编码写了col2
  • 做数据转换时字段被重命名,你没有同步修改访问代码
排查步骤
  • 在你代码报错的第88行前加调试打印,输出当前Row的所有字段:print(row.__fields__),直接确认字段列表里是否存在你要访问的col2
  • 打印row.asDict()输出当前行的全量键值,确认你要取的时间戳字段对应的真实字段名
  • 检查你调用foreachPartition/mapPartitions之前,DataFrame的select语句是否真的包含了两个删除条件对应的字段
修复方案
  1. 修正字段访问逻辑,不要硬编码写错的字段名,优先用字典下标方式访问Row字段,比属性访问更不容易踩拼写坑
  2. 传给psycopg2的参数必须转成原生Python类型的元组,不要直接传Spark Row对象,避免类型兼容问题
  3. 数据库连接必须在每个Executor分区内部初始化,不要在Driver端创建连接后序列化传到Executor,会触发连接失效错误

参考修复后的分区处理代码:

def execute_batch_delete(partition_iter):
    # 每个分区内独立初始化数据库连接
    conn = psycopg2.connect(
        host="your_aurora_endpoint",
        port=5432,
        dbname="your_db",
        user="your_user",
        password="your_pwd"
    )
    cur = conn.cursor()
    batch_data = []
    for row in partition_iter:
        # 替换成你打印出来的真实字段名,不要硬写col2
        val1 = row["real_col1_name"]
        val2 = row["real_col2_name"]
        batch_data.append( (val1, val2) )
        # 凑够批次量执行删除
        if len(batch_data) >= 2000:
            delete_sql = "DELETE FROM {target_table} WHERE (col1, col2) IN ( VALUES %s)".format(target_table=table)
            extras.execute_values(cur, delete_sql, batch_data, template="(%s, %s)", page_size=2000)
            conn.commit()
            batch_data.clear()
    # 处理最后一批剩余数据
    if batch_data:
        delete_sql = "DELETE FROM {target_table} WHERE (col1, col2) IN ( VALUES %s)".format(target_table=table)
        extras.execute_values(cur, delete_sql, batch_data, template="(%s, %s)", page_size=2000)
        conn.commit()
    cur.close()
    conn.close()

# 调用前确保DataFrame选对了字段
# delete_source_df = source_df.select("real_col1_name", "real_col2_name")
delete_source_df.rdd.foreachPartition(execute_batch_delete)
补充说明

你当前配置的page_size=2000是合理的:PostgreSQL协议单条SQL的参数上限是65535,2000条数据对应4000个参数,远低于阈值,不会触发参数过多的报错。等你修复完字段访问的问题后,这段删除逻辑可以正常运行。

内容的提问来源于stack exchange,提问作者이병준

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 17:51:23