如何使用Python将AWS PostgreSQL查询结果写入S3存储桶
问题原因与修复方案
核心错误点
cursor.copy_expert()的SQL入参不符合要求:该方法只能接收PostgreSQL的COPY格式语句,不能直接传入普通SELECT查询,你之前的写法会触发SQL语法错误,只是因为异常逻辑只打印错误不终止程序,所以你误以为代码执行成功。- 冗余执行了
cursor.execute(copy_query):普通SELECT查询和COPY导出是两个独立操作,不需要先执行一次SELECT。 - 硬编码AK/SK存在安全风险:本地测试可以通过AWS CLI的凭证文件配置,Lambda环境直接绑定S3写入权限的IAM角色即可,无需写死密钥。
修正后完整代码
import psycopg2 import boto3 import os import io # 本地测试可保留,Lambda环境建议删除,自动读取角色权限 os.environ['AWS_ACCESS_KEY_ID'] = "XXXXXXXXXXXXX" os.environ['AWS_SECRET_ACCESS_KEY'] = "XXXXXXXXXXXX" endpoint = "rds_endpoint" username = 'user_name' etl_password = 'stored_pass' database_name = 'db_name' s3_resource = boto3.resource('s3') file_name = 'daily_export' bucket = "my s3 bucket" # 调整为COPY格式语句,指定导出为CSV、带表头 copy_query = ''' COPY ( select parent.brand as business_type , app.business_name as business_name from hdsn_rsp parent join apt_ds app on parent.id = app.id ) TO STDOUT WITH ( FORMAT csv, HEADER true, ENCODING 'utf8' ) ''' def export_postgres_to_s3(): connection = None try: connection = psycopg2.connect( user= username, password= etl_password, host= endpoint, port="5432", database= database_name ) cursor = connection.cursor() csv_buffer = io.StringIO() # 直接执行COPY语句写入内存缓冲区 cursor.copy_expert(copy_query, csv_buffer) # 上传到S3 s3_resource.Object(bucket, f"{file_name}.csv").put(Body=csv_buffer.getvalue()) print(f"文件上传成功,大小:{len(csv_buffer.getvalue())} 字节") return "文件已成功写入S3存储桶" except(Exception, psycopg2.Error) as error: print("执行失败:", error) # 出错时明确返回失败信息,不要误导执行结果 return f"执行失败:{str(error)}" finally: if connection: cursor.close() connection.close() # 本地测试直接调用 if __name__ == "__main__": res = export_postgres_to_s3() print(res)
验证注意事项
- 运行时注意查看控制台打印的错误信息,如果是PostgreSQL的语法错误,可根据提示调整COPY语句的参数。
- 如果要迁移到Lambda运行,直接把函数名改成
lambda_handler,删除硬编码的AK/SK,给Lambda执行角色绑定S3存储桶的s3:PutObject权限即可。
内容的提问来源于stack exchange,提问作者26Cocktails
相关产品推荐
相关产品推荐

