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

多线程模式下使用aiomysql失败的技术求助

问题分析与解决方案

首先得明确:aiomysql是基于asyncio实现的异步数据库库,而asyncio的事件循环本质是单线程运行的——这就是你遇到问题的核心原因。

你提到单线程下正常运行,是因为在单线程环境中,asyncio的事件循环处于活跃状态,能够调度并执行await cur.execute(...)这类异步操作;但到了多线程模式下,你大概率是在新线程里直接调用了异步代码,却没有为这个线程初始化并启动asyncio事件循环,导致await语句无法被事件循环调度执行,代码直接卡在这一步,自然不会走到failed here的打印。

具体解决思路

针对多线程场景使用aiomysql,有两种常见的正确姿势:

1. 为每个线程创建独立的事件循环

如果你的业务确实需要每个线程独立处理数据库操作,可以在新线程的入口函数中,手动创建并运行asyncio事件循环,把异步数据库操作包装在协程里交给循环执行:

import asyncio
import threading
import aiomysql

async def db_query():
    print('ok here')
    try:
        pool = await aiomysql.create_pool(host='localhost', user='root', password='xxx', db='test')
        async with pool.acquire() as conn:
            async with conn.cursor() as cur:
                await cur.execute('select * from test')
                result = await cur.fetchall()
                print(result)
    except Exception as e:
        print('failed here', e)
    finally:
        pool.close()
        await pool.wait_closed()

def thread_func():
    # 为当前线程创建并设置事件循环
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    # 运行异步协程
    loop.run_until_complete(db_query())
    loop.close()

# 启动多线程
threading.Thread(target=thread_func).start()

2. 将异步任务提交到主线程的事件循环

如果不需要每个线程独立维护连接池,更推荐的方式是在主线程启动事件循环,然后用asyncio.run_coroutine_threadsafe把数据库操作的协程提交到主线程的事件循环中执行:

import asyncio
import threading
import aiomysql

async def db_query():
    print('ok here')
    try:
        pool = await aiomysql.create_pool(host='localhost', user='root', password='xxx', db='test')
        async with pool.acquire() as conn:
            async with conn.cursor() as cur:
                await cur.execute('select * from test')
                result = await cur.fetchall()
                print(result)
    except Exception as e:
        print('failed here', e)
    finally:
        pool.close()
        await pool.wait_closed()

def thread_func(main_loop):
    # 将协程提交到主线程的事件循环
    future = asyncio.run_coroutine_threadsafe(db_query(), main_loop)
    # 等待任务完成(可选,根据业务需求)
    future.result()

if __name__ == '__main__':
    # 主线程启动事件循环
    main_loop = asyncio.get_event_loop()
    # 启动线程,传入主线程的事件循环
    threading.Thread(target=thread_func, args=(main_loop,)).start()
    # 主线程保持事件循环运行
    main_loop.run_forever()

关键提醒

  • 不要在多线程中直接调用await语句,必须确保await所在的协程是在一个已启动的asyncio事件循环中执行的。
  • aiomysql的连接池是绑定到创建它的事件循环的,所以如果用多线程独立循环的方式,每个线程都需要创建自己的连接池;如果用主线程循环的方式,连接池可以复用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 08:23:29