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

使用Airflow、Python CSV及Psycopg优化数据库的拉取与推送操作

Airflow、Python CSV及Psycopg优化数据库的拉取与推送操作

兄弟,你的担心完全是对的——当前代码把全量Oracle查询结果都塞进内存里的io.StringIO,一旦数据量达到百万甚至千万级,内存直接就爆了,Airflow任务的资源配额本来就有限,这绝对是高风险操作。我给你拆解下现有问题,再给你最优的低内存优化方案:

一、现有代码的核心内存问题

你当前的create_query_buffer函数会把Oracle返回的所有数据全部写入内存中的StringIO缓存,才会返回给后续的PostgreSQL写入步骤。这意味着整个数据集都要占用内存空间——如果单表数据是10GB,你的任务内存就得至少预留10GB以上,完全不符合低内存的需求。哪怕你调小read(INT)的参数,也只是改变读取缓存的方式,缓存里的全量数据已经占了内存,根本解决不了根本问题。

二、最优低内存方案:流式处理(边读边写)

不需要把数据写到磁盘再读(磁盘IO反而会拖慢速度,还占磁盘空间),最优方案是流式处理:一边从Oracle批量拉取数据,一边直接写入PostgreSQL的copy操作中,内存里永远只保留当前批次的行数据,不会存全量数据集。

三、优化后的代码实现

1. 重构数据拉取逻辑,去掉全量缓存

把原来的create_query_buffer改成直接流式写入PostgreSQL的copy对象的函数,不再返回全量StringIO:

def stream_oracle_to_postgres_copy(oracle_con, sql, pg_copy_obj):
    import csv
    import io
    import time

    cur = oracle_con.cursor()
    # 保持你原来的游标优化参数,按需调整
    cur.prefetchrows = 17000
    cur.arraysize = 20000
    cur.execute(sql)
    start_time = time.time()

    # 先写入列名(CSV首行)
    col_names = [row[0] for row in cur.description]
    col_buffer = io.StringIO()
    csv.writer(col_buffer, lineterminator="\n").writerow(col_names)
    pg_copy_obj.write(col_buffer.getvalue())
    col_buffer.close()

    # 批量拉取并流式写入
    while True:
        rows = cur.fetchmany()
        if not rows:
            break
        # 将当前批次行转成CSV格式字符串,写入PostgreSQL copy
        batch_buffer = io.StringIO()
        csv.writer(batch_buffer, lineterminator="\n").writerows(rows)
        pg_copy_obj.write(batch_buffer.getvalue())
        batch_buffer.close()

    print(
        f"Time to stream oracle data to PostgreSQL: {round(time.time() - start_time,2)} seconds"
    )
    cur.close()

2. 调整主任务逻辑,结合流式处理

修改stage_ods任务,直接调用上面的流式函数,去掉中间的全量缓存环节:

from airflow.decorators import task
import oracledb
import psycopg
from psycopg import sql

@task
def stage_ods(
    queries: list[str], table_names: list[str], postgres_con: str, oracle_con: str
) -> None:
    oracle_conn = oracledb.connect(oracle_con)

    with psycopg.connect(postgres_con, autocommit=True) as pg_conn:
        pg_cur = pg_conn.cursor()

        for query, table in zip(queries, table_names):
            # 先获取列名(用于构造COPY语句)
            oracle_cur = oracle_conn.cursor()
            oracle_cur.execute(query)
            col_names = [col[0].lower() for col in oracle_cur.description]
            oracle_cur.close()

            # 截断目标表
            with pg_conn.transaction():
                pg_cur.execute(
                    sql.SQL("TRUNCATE TABLE {} RESTART IDENTITY CASCADE").format(
                        sql.Identifier(table)
                    )
                )

                # 启动PostgreSQL COPY操作,流式写入
                with pg_cur.copy(
                    sql.SQL("COPY {} ({}) FROM STDIN WITH CSV").format(
                        sql.Identifier(table),
                        sql.SQL(", ").join(map(sql.Identifier, col_names)),
                    )
                ) as copy_obj:
                    stream_oracle_to_postgres_copy(oracle_conn, query, copy_obj)

    oracle_conn.close()

if __name__ == "__main__":
    stage_ods()

四、关键优化点解析

  1. 彻底消除全量内存缓存:内存中只会保留当前fetchmany批次的行数据,以及临时的小批量CSV缓存,内存占用被严格控制在单批次数据大小范围内。
  2. 流式传输效率更高:避免了先写全量缓存再读的中间环节,数据从Oracle直接流向PostgreSQL,速度更快。
  3. 保留原有游标优化:继续使用prefetchrows和arraysize参数优化Oracle批量拉取的效率,平衡单批次内存占用和拉取速度。

五、关于“写磁盘还是内存缓存”的疑问

完全不需要写到磁盘:

  • 磁盘IO的速度远慢于内存操作,会大幅增加任务执行时间;
  • Airflow任务的临时磁盘空间可能有限,大表数据写入磁盘可能触发空间不足;
  • 流式处理的内存占用已经足够低,比写磁盘的方案更高效、更省资源。

额外注意事项

  • 调整游标参数:如果单条数据行很大(比如包含长文本),可以适当减小arraysize,避免单批次内存占用过高;
  • 资源配置:给Airflow的这个task设置合理的内存配额(比如1-2GB,根据单批次大小调整),避免资源浪费或不足;
  • 异常处理:可以在代码中加入try-except-finally块,确保连接和游标在异常时能正确关闭,避免资源泄漏。

备注:内容来源于stack exchange,提问作者moth

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 17:53:11