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

如何优化Python抽取Oracle大表并上传至S3的性能?

问题描述

我用Python脚本连接Oracle数据库抽取10张表的数据,其中一张3GB的表抽取并上传至S3耗时约4小时。请问如何优化以下Python脚本的性能?使用Parquet等非CSV格式能否提升性能?恳请提供建议或解决方案。

现有代码:

def extract_handler():
    # Parameters defined in cloudwatch event
    env = os.environ['Environment'] if 'Environment' in os.environ else 'sit'

    # FTP parameters
    host = f"/{env}/connet_HOSTNAME"
    username = f"/{env}/connect_USERNAME"
    password = f"/{env}/connect_PASSWORD"

    host = get_parameters(host)
    username = get_parameters(username)
    password = get_parameters(password)
    today = date.today()
    current_date = today.strftime("%Y%m%d")
    con = None
    cur = None
    tables = ["table1", "table2","table3"........."table10"]
    bucket = "bucket_name"

    for table in tables:
        try:
            con = cx_Oracle.connect(username, password, host, encoding="UTF-8")
            cur = con.cursor()
            logging.info('Successfully established the connection to Oracle db')
            table_name = table.split(".")[1]
            logging.info("######## Table name:"+ table+" ###### ")
            logging.info("****** PROCESSING:" +table_name+" *********")
            cur.execute("SELECT count(*) FROM  {}".format(table))
            count = cur.fetchone()[0]
            logging.info("Count:", count)
            if count > 0:
                cur1 = con.cursor()
                # Define the desired timestamp format
                timestamp_format = '%Y/%m/%d %H:%M:%S'
                # Execute a query to read a table
                cur1.execute( "select * from {} where TRUNC(DWH_CREATED_ON)=TRUNC(SYSDATE)-1".format(table))
                batch_size = 10000
                rows = cur1.fetchmany(batch_size)
                csv_file = f"/tmp/{table_name}.csv"
                with open(csv_file, "w", newline="") as f:
                    # Add file_date column as the first column
                    writer = csv.DictWriter(f, fieldnames=['file_date'] + [col[0] for col in cur1.description],
                                            delimiter='\t')
                    writer.writeheader()
                    logging.info("Header added to the table:" + table + "######")
                while rows:
                    for row in rows:
                        row_dict = {'file_date': current_date}
                        for i, col in enumerate(cur1.description):
                            if col[1] == cx_Oracle.DATETIME:
                                if row[i] is not None:
                                    row_dict[col[0]] = row[i].strftime(timestamp_format)
                                else:
                                    row_dict[col[0]] = ""
                            else:
                                row_dict[col[0]] = row[i]

                        with open(csv_file, "a", newline="") as f:
                            # Add file_date column as the first column
                            writer = csv.DictWriter(f, fieldnames=['file_date'] + [col[0] for col in cur1.description],
                                                    delimiter='\t')
                            writer.writerow(row_dict)
                        # Fetch the next batch of 100 rows
                        rows = cur1.fetchmany(batch_size)
                logging.info("Records written to the temp file for the table :" + table + "######")
                s3_path = "NorthernRegion" + '/' + table_name + '/' + current_date + '/' + table_name + '.csv'
                s3_client = boto3.client('s3', region_name='region-central-1')
                s3_client.upload_file('/tmp/' + table_name + '.csv', bucket, s3_path)
                logging.info(table + "File uploaded to S3 ######")
            else:
                logging.info('Table not having data')
                return 'Data is not refreshed yet, Hence quitting..'
            if cur1:
                cur1.close()

        except Exception as err:
            #Handle or log other exceptions such as bucket doesn't exist
            logging.error(err)
        finally:
            if cur:
                cur.close()
            if con:
                con.close()
        return "Successfully processed"

核心性能问题分析

现有代码的低效点集中在:

  • 逐行文件IO:每写一条数据就打开/关闭一次文件,IO开销极大
  • 单条数据处理:逐行转换字典、写入,CPU利用率极低
  • 重复资源初始化:每张表都重新建立Oracle连接、创建S3客户端,浪费连接开销
  • 冗余查询:单独执行count(*)会触发额外全表扫描
  • CSV格式缺陷:文本存储体积大、读写慢,无压缩,导致磁盘和网络传输耗时久

优化方案及Parquet格式的作用

1. 改用Parquet格式(强烈推荐)

Parquet是列式存储格式,相比CSV能带来显著性能提升:

  • 体积压缩:通常可将数据压缩至原CSV的1/5~1/10,大幅减少磁盘IO和S3上传时间
  • 读写效率:列式存储适配批量处理逻辑,Python的pyarrow/pandas对其支持高效,避免逐行字符串转换的开销
  • 类型保留:自动保留日期、数值等原生类型,无需手动格式化转换

2. 批量处理优化

  • 保持文件全程打开状态,避免重复打开/关闭的IO开销
  • 积累批量数据后一次性写入,而非逐行处理
  • 增大数据库批量读取的batch_size(如调至10万,根据内存情况调整)

3. 数据库操作优化

  • 仅建立一次Oracle连接,所有表共用该连接,消除重复连接的开销
  • 移除SELECT count(*)查询,直接通过游标迭代判断是否有数据,避免额外全表扫描
  • 优化SQL查询:将TRUNC(DWH_CREATED_ON)=TRUNC(SYSDATE)-1改为范围查询DWH_CREATED_ON >= TRUNC(SYSDATE)-1 AND DWH_CREATED_ON < TRUNC(SYSDATE),让Oracle能利用字段索引加速查询

4. S3上传优化

  • 仅初始化一次S3客户端,复用连接
  • 大文件上传可启用多线程(boto3默认支持,也可显式配置)
  • 即使保留CSV格式,也建议启用gzip压缩后再上传,减少传输体积

优化后的代码示例

import os
import logging
import cx_Oracle
import boto3
import pandas as pd
from datetime import date

def extract_handler():
    # 参数获取
    env = os.environ.get('Environment', 'sit')
    host = get_parameters(f"/{env}/connet_HOSTNAME")
    username = get_parameters(f"/{env}/connect_USERNAME")
    password = get_parameters(f"/{env}/connect_PASSWORD")
    current_date = date.today().strftime("%Y%m%d")
    tables = ["table1", "table2", "table3", ..., "table10"]
    bucket = "bucket_name"

    # 初始化全局资源
    con = None
    s3_client = boto3.client('s3', region_name='region-central-1')

    try:
        # 仅建立一次Oracle连接
        con = cx_Oracle.connect(username, password, host, encoding="UTF-8")
        logging.info('Oracle数据库连接成功')

        for table in tables:
            logging.info(f"######## 开始处理表: {table} ###### ")
            table_name = table.split(".")[1]
            
            # 优化后的SQL查询,利用索引加速
            query = """
                SELECT * 
                FROM {} 
                WHERE DWH_CREATED_ON >= TRUNC(SYSDATE)-1 
                  AND DWH_CREATED_ON < TRUNC(SYSDATE)
            """.format(table)
            
            # 用pandas批量读取数据,chunksize根据内存调整
            df_iter = pd.read_sql(query, con, chunksize=100000)
            
            parquet_file = f"/tmp/{table_name}.parquet"
            first_chunk = True
            
            for df in df_iter:
                # 添加file_date列并调整顺序
                df['file_date'] = current_date
                cols = ['file_date'] + [col for col in df.columns if col != 'file_date']
                df = df[cols]
                
                # 批量写入Parquet,第一次创建文件,后续追加
                df.to_parquet(
                    parquet_file,
                    engine='pyarrow',
                    append=not first_chunk,
                    compression='snappy'  # 平衡压缩率和速度,可选gzip获得更高压缩率
                )
                first_chunk = False
                logging.info(f"已写入{len(df)}条数据到临时文件: {table_name}")
            
            # 检查是否有数据并上传
            if not first_chunk:
                logging.info(f"表{table}数据写入临时文件完成")
                s3_path = f"NorthernRegion/{table_name}/{current_date}/{table_name}.parquet"
                s3_client.upload_file(parquet_file, bucket, s3_path)
                logging.info(f"表{table}文件已上传至S3: {s3_path}")
            else:
                logging.info(f"表{table}无数据")
        
        return "所有表处理完成"

    except Exception as err:
        logging.error(f"处理出错: {str(err)}")
        raise
    finally:
        # 统一关闭资源
        if con:
            con.close()
            logging.info("Oracle连接已关闭")

# 假设get_parameters为已实现的参数获取函数(如从SSM Parameter Store读取)
def get_parameters(param_path):
    # 此处实现参数获取逻辑
    pass

额外优化建议

  • 内存适配:如果运行环境内存有限,调小chunksize避免内存溢出
  • 并行处理:在Oracle连接并发允许的前提下,用多线程同时处理不同表
  • 驱动更新:确保使用最新版cx_Oracle,利用官方底层性能优化
  • S3内网访问:在AWS环境内运行时,配置VPC端点访问S3,避免公网传输开销

内容的提问来源于stack exchange,提问作者marjun

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 23:00:27