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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 10:07:49