数据库与异步编程结合:解决并发循环中数据库处理延迟问题
数据库操作与循环的异步执行解决方案
你的问题核心是:同步的数据库操作阻塞了异步事件循环,导致多个协程无法真正并行,甚至数据库还没完成写入,其他逻辑就已经进入下一轮判断。以下是具体解决方法:
1. 替换为异步数据库驱动(最优方案)
如果你的数据库支持异步驱动(比如PostgreSQL用asyncpg、MySQL用aiomysql),直接把同步的数据库方法改成异步调用,确保数据库操作不会阻塞事件循环:
async def send_mes_to_users(client): async with client: user = 'userrrrrrrrr' try: # 异步检查用户是否已邀请 if await db.check_invited(name_table=table, user_name=user) != "TRUE": await client.send_message(entity=user, message=message) # 异步写入数据库 await db.add_invited(name_table=table, user=user) print(f'Sent: @{user}') except errors.BadRequestError as e: print(f'Error: {e}') async def main(): await asyncio.gather( send_mes_to_users(TelegramClient('user1', api_id, api_hash)), send_mes_to_users(TelegramClient('user2', api_id, api_hash)), ) asyncio.run(main())
2. 用线程池包装同步数据库操作
如果没有异步驱动可用,用asyncio.to_thread把同步的数据库操作放到线程池执行,避免阻塞异步事件循环:
async def send_mes_to_users(client): async with client: user = 'userrrrrrrrr' try: # 在线程池里执行同步检查 check_result = await asyncio.to_thread( db.check_invited, name_table=table, user_name=user ) if check_result != "TRUE": await client.send_message(entity=user, message=message) # 在线程池里执行同步写入 await asyncio.to_thread( db.add_invited, name_table=table, user=user ) print(f'Sent: @{user}') except errors.BadRequestError as e: print(f'Error: {e}')
3. 处理竞态条件
即使改成异步,多个协程同时操作数据库可能出现同一用户被重复处理的情况,可通过两种方式解决:
- 数据库层面:给
user_name字段添加唯一索引,避免重复写入。 - 代码层面:用
asyncio.Lock加锁,确保同一时间只有一个协程执行数据库检查和写入:
# 定义全局锁(单进程场景下使用) db_operation_lock = asyncio.Lock() async def send_mes_to_users(client): async with client: user = 'userrrrrrrrr' try: async with db_operation_lock: check_result = await db.check_invited(name_table=table, user_name=user) if check_result != "TRUE": await client.send_message(entity=user, message=message) await db.add_invited(name_table=table, user=user) print(f'Sent: @{user}') except errors.BadRequestError as e: print(f'Error: {e}')
内容的提问来源于stack exchange,提问作者Сергей Викторович
相关产品推荐
相关产品推荐

