使用Python将S3数据加载至AWS RDS Postgres的标准实现方案
实现方案指导
前置准备
首先要在你的RDS PostgreSQL实例上完成基础配置:
- 给RDS实例关联具备S3桶读取权限的IAM角色,确保RDS实例可以访问目标S3资源
- 连接到PostgreSQL实例,执行以下SQL开启扩展:
CREATE EXTENSION IF NOT EXISTS aws_s3 CASCADE; -- 给你使用的数据库用户授予扩展使用权限 GRANT USAGE ON SCHEMA aws_s3 TO your_pg_user;
推荐实现方案(直接通过扩展从S3拉取数据,性能最优)
这个方案和你现有GCP逻辑最接近,不需要把文件下载到本地,直接由PostgreSQL从S3拉取数据写入表,适合大文件场景,Python侧只需要调用扩展提供的导入函数即可。
依赖安装
需要先安装Python的PostgreSQL连接库:
pip install psycopg2-binary
代码示例
import psycopg2 from psycopg2 import sql # 配置参数 PG_CONFIG = { "host": "你的RDS实例地址", "port": 5432, "user": "你的PG用户名", "password": "你的PG密码", "dbname": "目标数据库名" } S3_BUCKET = "你的S3桶名" S3_FILE_PATH = "path-to-your/file.json" AWS_REGION = "你的S3桶所在区域,比如us-east-1" RDS_IAM_ROLE_ARN = "你关联到RDS的IAM角色ARN" TARGET_TABLE = "你的目标表名,格式为schema.table" # COPY命令参数,对应你要导入的ndjson格式,按需调整 COPY_OPTIONS = "FORMAT json, DELIMITER '\n'" # 建立PG连接 conn = psycopg2.connect(**PG_CONFIG) cur = conn.cursor() try: # 调用aws_s3扩展的导入函数 import_query = sql.SQL(""" SELECT aws_s3.table_import_from_s3( %s, '', -- 留空表示导入所有列,也可指定列名,比如 'col1,col2,col3' %s, aws_commons.create_s3_uri(%s, %s, %s), %s ); """) cur.execute(import_query, ( TARGET_TABLE, COPY_OPTIONS, S3_BUCKET, S3_FILE_PATH, AWS_REGION, RDS_IAM_ROLE_ARN )) conn.commit() # 查询导入的行数 cur.execute(sql.SQL("SELECT COUNT(*) FROM {};").format(sql.Identifier(*TARGET_TABLE.split('.')))) total_rows = cur.fetchone()[0] print(f"导入完成,当前表总行数:{total_rows}") except Exception as e: conn.rollback() print(f"导入失败:{str(e)}") raise finally: cur.close() conn.close()
备选实现方案(无需配置扩展,适合小数据量场景)
如果你不想配置扩展,也可以先通过boto3把S3文件读取到本地内存,再写入PG:
依赖安装
pip install boto3 pandas psycopg2-binary sqlalchemy
代码示例
import boto3 import pandas as pd from sqlalchemy import create_engine # 初始化S3客户端 s3 = boto3.client('s3') # 初始化PG引擎 pg_engine = create_engine("postgresql://用户名:密码@RDS地址:5432/数据库名") # 读取S3上的ndjson文件 response = s3.get_object(Bucket="你的桶名", Key="文件路径") df = pd.read_json(response['Body'], lines=True) # 写入PG,if_exists设为append对应追加写入,也可设为replace覆盖 df.to_sql( name="目标表名", con=pg_engine, schema="目标schema名", if_exists="append", index=False ) print(f"成功导入{len(df)}行数据")
Airflow适配说明
你可以直接把上述代码封装成Airflow的PythonOperator运行,也可以用Airflow自带的PostgresOperator直接执行导入SQL,不需要额外依赖工具,完全符合你自有架构的要求。
注意事项
- 如果你选择追加写入的模式,建议先把数据导入临时表,做去重/清洗后再合并到正式表,避免重复导入
- ndjson格式导入需要确保每一行的JSON结构和PG表的字段类型对应,如果有嵌套结构可以先在PG表中定义JSON类型的字段存储
- 大文件导入优先选择第一种扩展方案,避免占用Airflow worker的内存和带宽
内容的提问来源于stack exchange,提问作者Canovice
相关产品推荐
相关产品推荐

