Python多线程操作MySQL数据库时出现连接丢失问题求助
多线程下MySQL连接丢失问题排查与修复
问题核心
你的代码在单线程操作MySQL时正常,但启用多线程后出现「Lost connection to MySQL server」错误,根源在于全局共享的MySQL连接和游标线程不安全,再加上线程启动逻辑错误、SQL语句存在注入风险,共同导致了连接异常。
错误原因分析
- 连接/游标线程不安全:
mysql.connector的连接和游标对象不支持多线程并发操作,多个线程同时调用同一个游标执行SQL、提交事务,会直接打乱连接状态,触发连接丢失。 - 线程启动逻辑错误:
Thread(...).start()返回None,你后续调用的thread.join()实际是对None执行,根本没起到等待线程的作用,所有线程瞬间并发,进一步加剧了连接冲突。 - SQL注入风险:直接拼接字符串到SQL语句中,不仅会引发安全问题,还可能因格式错误导致SQL执行异常,间接触发连接问题。
修复方案
1. 给每个线程分配独立的MySQL连接
将连接和游标创建移到线程函数内部,确保每个线程拥有专属的连接资源,避免共享冲突。
修改book函数及依赖方法:
def book(user): # 每个线程单独初始化连接和游标 mydb = mysql.connector.connect( host="localhost", user="********", password="********", database="airline_checkin" ) mycursor = mydb.cursor() print(f"{user[1]} with id {user[0]} trying to book seat ") seats = available_seats(mycursor) print(f"Total Available seats: {seats[0]}") # 随机选择座位 seat_id = random.randint(1, seats[0]) seat = GetSeat(seat_id, mycursor) print(f"Seat selected: {seat[1]}") # 使用参数化查询避免SQL注入 update_sql = "UPDATE seats SET USER_ID = %s WHERE id = %s" mycursor.execute(update_sql, (user[0], seat[0])) mydb.commit() # 关闭当前线程的连接资源 mycursor.close() mydb.close()
同时调整依赖函数,让它们接收游标作为参数:
def available_seats(cursor): cursor.execute("SELECT COUNT(*) FROM seats WHERE USER_ID = 0") myresult = cursor.fetchall() return myresult[0] def GetUser(userID, cursor): cursor.execute("SELECT * FROM users WHERE id = %s", (userID,)) myresult = cursor.fetchall() return myresult[0] if myresult else None def GetSeat(seatID, cursor): cursor.execute("SELECT * FROM seats WHERE id = %s", (seatID,)) myresult = mycursor.fetchall() return myresult[0] if myresult else None
2. 修正线程启动与等待逻辑
先创建线程对象并保存,启动后统一等待所有线程执行完成:
# 驱动代码修改 reset() threads = [] for i in range(120): user = GetUser(i, mycursor) # 主线程中获取用户信息 if user: thread = Thread(target=book, args=(user,)) threads.append(thread) thread.start() # 等待所有线程执行完毕 for thread in threads: thread.join() input("Press enter to continue...") print() print("Here is the updated seating chart") airline()
3. 全局操作函数保留主线程连接
reset和airline函数在单线程环境下执行,可继续使用主线程的连接:
def reset(): mycursor.execute("UPDATE seats SET USER_ID = 0") mydb.commit() def airline(): seats = [[". " for i in range(20)] for j in range(7)] mycursor.execute("SELECT * FROM seats WHERE USER_ID > 0") myresult = mycursor.fetchall() for x in myresult: nm = x[1] row, c = nm.split("-") col = ALPHA_NUM[c] seats[col][int(row)-1] = "X " for i in range(0, int(len(seats)/2)): for j in range(len(seats[i])): print(seats[i][j], end="") print() for i in range(2): for j in range(len(seats[i])-1): print(" ", end="") print() for i in range(int(len(seats)/2), len(seats) - 1): for j in range(len(seats[i])): print(seats[i][j], end="") print()
可选优化:使用连接池减少开销
如果担心频繁创建连接的性能损耗,可以用连接池复用资源,同时保证线程安全:
from mysql.connector import pooling # 初始化连接池 db_pool = pooling.MySQLConnectionPool( pool_name="airline_pool", pool_size=10, # 根据并发量调整 host="localhost", user="********", password="********", database="airline_checkin" ) def book(user): # 从连接池获取连接 mydb = db_pool.get_connection() mycursor = mydb.cursor() # ... 业务逻辑和之前一致 ... mycursor.close() mydb.close() # 将连接归还到池
内容的提问来源于stack exchange,提问作者Suraj Singh
相关产品推荐
相关产品推荐

