Python脚本优化:大数据集下PostgreSQL插入与AWS SQS推送提速方案
优化方案与实现模式
1. 纠正认知:Python多线程完全适用于你的场景
Python的GIL对IO密集型任务(数据库插入、SQS推送这类涉及网络/磁盘等待的操作)几乎没有影响——IO等待时GIL会自动释放,所以多线程可以有效并行处理这两个任务,大幅提升整体效率。
2. 核心优化:批量操作(最见效的手段)
你的性能瓶颈主要在SQS单条推送的网络开销,加上单条数据库插入的累积耗时,先从批量操作入手:
PostgreSQL批量插入
用psycopg2的executemany或者更高效的COPY FROM(适合超大数据量):
import psycopg2 from psycopg2 import sql from io import StringIO # 方案1:executemany批量插入(适合中等数据量) def batch_insert_to_db(records): conn = psycopg2.connect("dbname=your_db user=your_user") cur = conn.cursor() insert_query = sql.SQL("INSERT INTO your_table (col1, col2) VALUES (%s, %s)") cur.executemany(insert_query, records) conn.commit() cur.close() conn.close() # 方案2:COPY FROM批量导入(适合1M级大数据量,速度比executemany快5-10倍) def copy_from_records(records): conn = psycopg2.connect("dbname=your_db user=your_user") cur = conn.cursor() buffer = StringIO() for row in records: buffer.write(f"{row[0]},{row[1]}\n") buffer.seek(0) cur.copy_from(buffer, "your_table", columns=('col1', 'col2'), sep=',') conn.commit() cur.close() conn.close()
SQS批量推送
AWS SQS支持send_message_batch接口,一次最多推送10条消息,大幅减少API调用次数:
import boto3 sqs = boto3.client('sqs', region_name='your_region') queue_url = 'your_queue_url' def batch_send_to_sqs(records): # 拆分成每10条一组(SQS批量接口的上限) batches = [records[i:i+10] for i in range(0, len(records), 10)] for batch in batches: entries = [ { 'Id': str(idx), 'MessageBody': str(row) # 根据你的数据格式序列化 } for idx, row in enumerate(batch) ] sqs.send_message_batch(QueueUrl=queue_url, Entries=entries)
3. 多线程并行处理(让插入和推送同时跑)
用concurrent.futures.ThreadPoolExecutor实现批量任务的并行处理,或者用生产者-消费者模式分离读数据、插数据库、推SQS三个环节:
生产者-消费者模式(适合大数据量)
import csv from concurrent.futures import ThreadPoolExecutor import queue # 全局队列存储读取到的批量数据,限制大小避免内存溢出 data_queue = queue.Queue(maxsize=10) def read_csv_to_queue(batch_size=100): with open('data.csv', 'r') as f: reader = csv.reader(f) next(reader) # 跳过表头 batch = [] for row in reader: batch.append(row) if len(batch) == batch_size: data_queue.put(batch) batch = [] # 处理最后一批不足batch_size的数据 if batch: data_queue.put(batch) data_queue.put(None) # 标记数据读取结束 def db_worker(): while True: batch = data_queue.get() if batch is None: data_queue.put(None) # 传递结束信号给其他worker break copy_from_records(batch) # 用最快的COPY FROM批量插入 def sqs_worker(): while True: batch = data_queue.get() if batch is None: break batch_send_to_sqs(batch) # 启动三个线程:读数据、写数据库、推SQS,并行执行 with ThreadPoolExecutor(max_workers=3) as executor: executor.submit(read_csv_to_queue, batch_size=100) executor.submit(db_worker) executor.submit(sqs_worker)
4. 异步IO的正确实现(解决你之前的问题)
你之前用asyncio的问题是用了同步数据库驱动,导致异步任务被阻塞。改用异步数据库驱动(如asyncpg),让插入和推送都异步执行,边处理边执行,不会攒到最后:
import asyncio import csv import asyncpg import aioboto3 async def async_insert_to_db(conn, row): await conn.execute('INSERT INTO your_table (col1, col2) VALUES ($1, $2)', row[0], row[1]) async def async_send_to_sqs(sqs_client, queue_url, row): await sqs_client.send_message(QueueUrl=queue_url, MessageBody=str(row)) async def process_row(conn, sqs_client, queue_url, row): # 并行执行两个异步任务 await asyncio.gather( async_insert_to_db(conn, row), async_send_to_sqs(sqs_client, queue_url, row) ) async def main(): # 建立异步数据库连接 conn = await asyncpg.connect(user='your_user', dbname='your_db') # 建立异步SQS客户端 session = aioboto3.Session() async with session.client('sqs', region_name='your_region') as sqs_client: queue_url = 'your_queue_url' # 读取CSV并异步处理 with open('data.csv', 'r') as f: reader = csv.reader(f) next(reader) tasks = [] # 控制并发数,避免打满数据库或SQS的并发限制 semaphore = asyncio.Semaphore(50) async def bounded_process(row): async with semaphore: await process_row(conn, sqs_client, queue_url, row) for row in reader: tasks.append(bounded_process(row)) await asyncio.gather(*tasks) await conn.close() asyncio.run(main())
无服务器环境适配建议
如果是Lambda这类无服务器环境:
- 控制批量大小和并发数,避免超出内存/时间限制(比如Lambda最大超时15分钟,1M数据可以分批次处理,每次处理10k条)
- 用Lambda异步调用或Step Functions拆分任务,把大任务拆成多个小任务并行执行
- 优先用批量操作,减少IO次数,降低冷启动和网络开销
内容的提问来源于stack exchange,提问作者Dipanshu Chaubey
相关产品推荐
相关产品推荐

