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
相关产品推荐
相关产品推荐

