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语句是否真的包含了两个删除条件对应的字段
修复方案
- 修正字段访问逻辑,不要硬编码写错的字段名,优先用字典下标方式访问Row字段,比属性访问更不容易踩拼写坑
- 传给psycopg2的参数必须转成原生Python类型的元组,不要直接传Spark Row对象,避免类型兼容问题
- 数据库连接必须在每个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,提问作者이병준
相关产品推荐
相关产品推荐

