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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 05:54:02