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

Python Requests+MySQL+多进程/线程高效实现方案咨询

解决Python多线程/进程结合Requests+MySQL的并发性能问题

首先得戳破你遇到的两个核心问题的根源:

1. ThreadPool(线程池)的问题:共享MySQL连接导致阻塞

你用multiprocessing.dummy.ThreadPool时,所有线程属于同一个进程(所以PID相同是正常的,线程本来就共享进程PID),但你在多个线程里复用了同一个MySQL连接和cursor——这是致命错误!MySQL的连接对象并不是线程安全的,多个线程同时操作同一个连接会导致内部状态混乱、锁等待,最终表现为速度和单线程差不多,甚至卡死。

线程池的正确改造方案:每个线程独立创建MySQL连接

每个线程在执行任务前,自己初始化专属的MySQL连接和cursor,任务结束后关闭(或者用连接池复用)。这样彻底避免线程间的连接竞争。

示例代码:

from multiprocessing.dummy import ThreadPool
import requests
import mysql.connector

def init_thread_connection():
    # 每个线程初始化自己的连接
    return mysql.connector.connect(
        host="your_host",
        user="your_user",
        password="your_pass",
        database="your_db"
    )

def run(url):
    # 每个线程用自己的独立连接
    mydb = init_thread_connection()
    mycursor = mydb.cursor()
    
    try:
        # 原子性获取最后一条ID(后面会说更优的任务分配方式)
        find_id = "SELECT id FROM Items ORDER BY id DESC LIMIT 1"
        mycursor.execute(find_id)
        last_id = mycursor.fetchone()[0]
        
        # 复用requests会话提升请求效率(避免重复建立TCP连接)
        with requests.Session() as sess:
            data = sess.get(url).text
        
        # 处理数据并写入数据库
        insert_sql = "INSERT INTO TargetTable (content, source_id) VALUES (%s, %s)"
        mycursor.execute(insert_sql, (data, last_id + 1))
        mydb.commit()
    except Exception as e:
        print(f"处理URL {url}出错: {str(e)}")
        mydb.rollback()
    finally:
        # 关闭当前线程的连接资源
        mycursor.close()
        mydb.close()

if __name__ == "__main__":
    urls = ["https://example.com/1", "https://example.com/2", ...]
    pool = ThreadPool(10)
    results = pool.map(run, urls)
    pool.close()
    pool.join()

2. 多进程的问题:无原子性任务分配导致重复读取+锁冲突

用multiprocessing.Process时,每个进程有独立的MySQL连接,但你让每个worker自己去读最后一条ID,这会导致多个进程同时读到同一个ID(因为普通SELECT操作不是原子的),进而重复处理;同时并发写入如果没有合理的锁机制,会触发MySQL的行锁/表锁,导致性能骤降。

多进程的正确改造方案:原子化任务分配+进程专属连接

方案A:数据库层面原子获取待处理任务

用SELECT ... FOR UPDATE SKIP LOCKED(MySQL 8.0+支持)来原子性获取并锁定未处理的任务,避免多个进程抢同一条数据:

from multiprocessing import Pool
import requests
import mysql.connector

def init_process_connection():
    # 每个进程初始化自己的连接(放在全局变量方便worker调用)
    global mydb, mycursor
    mydb = mysql.connector.connect(
        host="your_host",
        user="your_user",
        password="your_pass",
        database="your_db"
    )
    mycursor = mydb.cursor()

def run(url):
    try:
        # 原子性获取并锁定一条待处理的任务(假设Items表有status字段标记状态)
        get_task_sql = """
            SELECT id FROM Items WHERE status = 'pending' 
            ORDER BY id DESC LIMIT 1 FOR UPDATE SKIP LOCKED
        """
        mycursor.execute(get_task_sql)
        task_id = mycursor.fetchone()
        if not task_id:
            return "无待处理任务"
        
        # 先标记任务为处理中,避免其他进程再读取
        update_status_sql = "UPDATE Items SET status = 'processing' WHERE id = %s"
        mycursor.execute(update_status_sql, (task_id[0],))
        mydb.commit()
        
        # 请求数据
        with requests.Session() as sess:
            data = sess.get(url).text
        
        # 写入结果并标记任务完成
        insert_sql = "INSERT INTO TargetTable (content, source_id) VALUES (%s, %s)"
        mycursor.execute(insert_sql, (data, task_id[0]))
        update_finish_sql = "UPDATE Items SET status = 'done' WHERE id = %s"
        mycursor.execute(update_finish_sql, (task_id[0],))
        mydb.commit()
    except Exception as e:
        print(f"处理URL {url}出错: {str(e)}")
        mydb.rollback()

if __name__ == "__main__":
    urls = ["https://example.com/1", "https://example.com/2", ...]
    # 初始化每个进程的数据库连接
    pool = Pool(10, initializer=init_process_connection)
    results = pool.map(run, urls)
    pool.close()
    pool.join()

方案B:主进程统一分发任务

由主进程一次性读取所有需要处理的任务,把URL和对应的ID绑定后放入队列,worker进程从队列取任务处理,彻底避免每个worker自己去数据库抢数据:

from multiprocessing import Pool, Queue
import requests
import mysql.connector

def worker(queue):
    # 每个worker进程初始化自己的连接和请求会话
    mydb = mysql.connector.connect(
        host="your_host",
        user="your_user",
        password="your_pass",
        database="your_db"
    )
    mycursor = mydb.cursor()
    sess = requests.Session()
    
    while True:
        task = queue.get()
        if task is None:  # 收到结束信号就退出
            break
        url, task_id = task
        try:
            data = sess.get(url).text
            insert_sql = "INSERT INTO TargetTable (content, source_id) VALUES (%s, %s)"
            mycursor.execute(insert_sql, (data, task_id))
            mydb.commit()
        except Exception as e:
            print(f"处理任务ID {task_id}出错: {str(e)}")
            mydb.rollback()
    
    # 关闭资源
    sess.close()
    mycursor.close()
    mydb.close()

if __name__ == "__main__":
    # 主进程读取所有待处理任务ID
    main_db = mysql.connector.connect(...)
    main_cursor = main_db.cursor()
    main_cursor.execute("SELECT id FROM Items ORDER BY id DESC")
    task_ids = [row[0] for row in main_cursor.fetchall()]
    main_cursor.close()
    main_db.close()
    
    # 绑定URL和任务ID(假设urls和task_ids一一对应)
    urls = ["https://example.com/1", "https://example.com/2", ...]
    tasks = list(zip(urls, task_ids))
    
    # 创建任务队列并放入所有任务
    queue = Queue()
    for task in tasks:
        queue.put(task)
    
    # 启动worker进程,每个进程最后发送一个None作为结束信号
    pool_size = 10
    pool = Pool(pool_size)
    for _ in range(pool_size):
        pool.apply_async(worker, args=(queue,))
    
    # 等待所有任务处理完成
    queue.join()
    # 发送结束信号
    for _ in range(pool_size):
        queue.put(None)
    pool.close()
    pool.join()

通用优化建议

  • 复用Requests会话:每个线程/进程用requests.Session(),可以复用TCP连接,大幅提升请求速度。
  • 数据库连接池:不用每次任务都创建/关闭连接,用mysql-connector的pool_size参数或者SQLAlchemy连接池来复用连接,减少连接开销。
  • 批量写入:如果业务允许,把多条数据攒到一起批量插入,减少commit次数,提升数据库写入效率。
  • 缩短事务时间:尽量把数据库操作的事务时间压缩到最短,减少锁等待的概率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 15:57:43