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

使用concurrent.futures多进程操作MySQL遇2003连接错误求助

多进程下AIOMySQL持续出现(2003, "Can't connect to MySQL server on '127.0.0.1'")问题排查与修复

问题现象

  • 脚本通过concurrent.futures创建多进程,所有进程向本地MySQL发送INSERT查询
  • 运行数小时后持续返回(2003, "Can't connect to MySQL server on '127.0.0.1'")错误,重试无效
  • MySQL服务处于活跃状态,MySQL Workbench可正常连接;重启MySQL/MariaDB服务无法解决,仅杀死脚本进程并重启脚本后恢复正常
  • 日志触发Connect函数中的#992分支

代码核心问题分析

从提供的代码来看,故障根源集中在连接管理混乱、异步与多进程兼容缺陷、错误处理逻辑漏洞三点:

  1. 连接泄漏严重

    • 每次请求都创建新连接,未使用连接池复用资源,高并发下会快速耗尽MySQL的max_connections上限,导致后续无法建立新连接
    • 异常分支中递归调用request_wrapper后,对conn的关闭逻辑存在漏洞:若连接未初始化则会报错,且递归创建的新连接未被妥善管理
    • conn.close()是同步方法,AIOMySQL连接需用await conn.close()才能真正释放,否则MySQL端会残留大量未关闭的连接
  2. 异步环境误用同步API

    • 代码中使用sleep(2)(同步阻塞)会冻结整个异步事件循环,导致异步任务堆积,连接资源无法及时释放,应替换为asyncio.sleep(2)
  3. 无限递归重试加剧资源消耗

    • Connect函数连接失败时无限制递归调用自身,当MySQL连接资源耗尽后,会持续占用CPU和内存,进一步恶化连接问题
  4. 多进程与异步资源冲突

    • concurrent.futures.ProcessPoolExecutor创建的子进程会复制父进程的事件循环状态,AIOMySQL连接绑定到事件循环,子进程复用父进程异步资源会导致异常,每个进程需独立初始化事件循环

修复方案

1. 改用AIOMySQL连接池

连接池可复用连接、限制最大连接数,避免频繁创建/销毁连接导致的资源耗尽,建议将maxsize设置为小于MySQL的max_connections值。

2. 替换同步sleep为异步版本

用asyncio.sleep(2)替代time.sleep(2),避免阻塞事件循环。

3. 限制重试次数,取消无限递归

在连接和请求逻辑中添加重试次数上限,超过次数后抛出异常或记录日志。

4. 多进程独立初始化异步环境

每个子进程单独创建事件循环和连接池,避免继承父进程的异步资源。

5. 用异步上下文管理器确保资源释放

使用async with管理连接和游标,自动处理关闭逻辑,避免手动关闭遗漏。

修复后的示例代码

import asyncio
import aiomysql
import datetime
from concurrent.futures import ProcessPoolExecutor

# 每个进程独立维护连接池
pool = None

async def init_pool(db='twitch_chat'):
    global pool
    pool = await aiomysql.create_pool(
        host=credentials.host,
        db=db,
        user=credentials.user,
        password=credentials.password,
        autocommit=True,
        maxsize=10,  # 根据MySQL max_connections调整
        minsize=2
    )

async def request_wrapper(query_type, request, dict=False, retry=0, charset='', collation=None):
    max_retry = 3
    try:
        async with pool.acquire() as conn:
            if query_type == 'INSERT' and len(charset) > 0:
                try:
                    await conn.set_charset_collation(charset=charset)
                except Exception as e:
                    print(e, "[#994]")
            async with conn.cursor() as cursor:
                await cursor.execute('SET NAMES utf8mb4')
                await cursor.execute("SET CHARACTER SET utf8mb4")
                await cursor.execute("SET character_set_connection=utf8mb4")
                await cursor.execute(request)
                if query_type == "SELECT":
                    return await cursor.fetchall()
    except Exception as e:
        try:
            if e[0] == 1062:
                return
        except:
            print("======================")
            print(e, "#00.5")
            print("======================")
            return
        if retry < max_retry:
            print("======================")
            print(f"({datetime.datetime.now().strftime('%H:%M')}) [MYSQL ERROR]", request)
            print(e)
            print("RETRY IN 2 SECONDS.")
            print("======================")
            await asyncio.sleep(2)
            return await request_wrapper(query_type, request, dict, retry+1, charset, collation)
        else:
            print(f"({datetime.datetime.now().strftime('%H:%M')}) [MYSQL ERROR] 重试次数耗尽", request)
            raise e

async def process_task(query_type, request, **kwargs):
    # 子进程初始化连接池
    await init_pool()
    result = await request_wrapper(query_type, request, **kwargs)
    # 任务结束后关闭连接池
    pool.close()
    await pool.wait_closed()
    return result

def run_async_task(query_type, request, **kwargs):
    # 子进程独立创建事件循环
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    result = loop.run_until_complete(process_task(query_type, request, **kwargs))
    loop.close()
    return result

# 多进程调用示例
if __name__ == "__main__":
    with ProcessPoolExecutor(max_workers=4) as executor:
        # 批量提交INSERT任务
        futures = [
            executor.submit(run_async_task, 'INSERT', "INSERT INTO your_table (col1) VALUES ('value1')")
            for _ in range(10)
        ]
        # 处理任务结果
        for future in futures:
            try:
                future.result()
            except Exception as e:
                print(f"任务执行失败: {e}")

额外优化建议

  • 检查MySQL的max_connections配置,确保连接池maxsize加上其他业务连接(如Workbench、监控工具)不超过该值
  • 开启MySQL慢查询日志,用SHOW PROCESSLIST定期查看连接状态,确认是否有大量睡眠或阻塞连接
  • 避免在多进程间共享异步资源,每个进程需独立初始化自身的事件循环和连接池

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 07:55:29