使用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()
四、关键优化点解析
- 彻底消除全量内存缓存:内存中只会保留当前
fetchmany批次的行数据,以及临时的小批量CSV缓存,内存占用被严格控制在单批次数据大小范围内。 - 流式传输效率更高:避免了先写全量缓存再读的中间环节,数据从Oracle直接流向PostgreSQL,速度更快。
- 保留原有游标优化:继续使用
prefetchrows和arraysize参数优化Oracle批量拉取的效率,平衡单批次内存占用和拉取速度。
五、关于“写磁盘还是内存缓存”的疑问
完全不需要写到磁盘:
- 磁盘IO的速度远慢于内存操作,会大幅增加任务执行时间;
- Airflow任务的临时磁盘空间可能有限,大表数据写入磁盘可能触发空间不足;
- 流式处理的内存占用已经足够低,比写磁盘的方案更高效、更省资源。
额外注意事项
- 调整游标参数:如果单条数据行很大(比如包含长文本),可以适当减小
arraysize,避免单批次内存占用过高; - 资源配置:给Airflow的这个task设置合理的内存配额(比如1-2GB,根据单批次大小调整),避免资源浪费或不足;
- 异常处理:可以在代码中加入
try-except-finally块,确保连接和游标在异常时能正确关闭,避免资源泄漏。
备注:内容来源于stack exchange,提问作者moth
相关产品推荐
相关产品推荐

