Python JobQueue重复发送消息问题:快速操作触发,调试无异常
问题分析与解决方案:频繁触发函数导致JobQueue重复发送消息
问题现象
当用户快速、频繁触发first函数时,即便user_last_len的值没有变化,JobQueue仍会重复发送同一条消息2次;但正常慢速操作时,消息仅发送1次,符合预期。调试过程中无法复现该问题,消息仅发送1次。
相关代码
user_last_len = {} async def first(update: Update, context: CallbackContext): async def check(context: CallbackContext): await asyncio.sleep(2) try: conn = connection_pool.get_connection() cursor = conn.cursor() user_id = update.message.from_user.id cursor.execute("SELECT LikeID FROM Likes WHERE LikedID = %s", (user_id,)) likes = cursor.fetchall() len_likes = len(likes) last_len = user_last_len.get(user_id, 0) if len_likes > 0 and last_len != len_likes: await update.message.reply_text(f"{len_likes}", reply_markup=show_markup) user_last_len[user_id] = len_likes context.job.data = True cursor.close() conn.close() except Exception as e: pass finally: await asyncio.sleep(1) job = context.job_queue.run_repeating(check, interval=300, first=0, data=None, name="dgc") await job.run(context.application)
问题根源
每次调用first函数都会创建并启动一个新的重复任务。当用户快速操作时,短时间内多次进入first,会生成多个独立的check任务并行执行:
- 这些任务几乎同时读取数据库数据,此时
user_last_len还没被更新 - 多个任务都通过了
last_len != len_likes的判断,进而重复发送消息 - 调试时操作速度慢,前一个任务已经完成
user_last_len的更新,所以不会触发重复发送
另外,全局字典user_last_len没有并发保护,多任务同时读写时存在竞态条件,进一步加剧了重复发送的概率。
解决方案
1. 避免重复创建任务
在创建新任务前,先检查并移除已存在的同名任务,确保同一时间只有一个check任务运行。
2. 保护全局字典的并发访问
使用锁确保对user_last_len的读写操作是原子性的,避免多任务同时操作导致的判断错误。
修改后的代码
import asyncio from threading import Lock user_last_len = {} user_lock = Lock() # 用于保护user_last_len的线程安全锁 async def first(update: Update, context: CallbackContext): # 先清除已存在的同名任务,避免重复执行 existing_jobs = context.job_queue.get_jobs_by_name("dgc") if existing_jobs: for job in existing_jobs: job.schedule_removal() async def check(context: CallbackContext): await asyncio.sleep(2) try: conn = connection_pool.get_connection() cursor = conn.cursor() user_id = update.message.from_user.id cursor.execute("SELECT LikeID FROM Likes WHERE LikedID = %s", (user_id,)) likes = cursor.fetchall() len_likes = len(likes) # 加锁确保读写user_last_len的原子性 with user_lock: last_len = user_last_len.get(user_id, 0) if len_likes > 0 and last_len != len_likes: await update.message.reply_text(f"{len_likes}", reply_markup=show_markup) user_last_len[user_id] = len_likes context.job.data = True cursor.close() conn.close() except Exception as e: pass finally: await asyncio.sleep(1) job = context.job_queue.run_repeating(check, interval=300, first=0, data=None, name="dgc") await job.run(context.application)
关键修改说明
- 任务去重:通过
get_jobs_by_name获取已有任务并移除,确保每次进入first函数时只有一个新任务启动。 - 并发保护:使用
threading.Lock包裹对user_last_len的读写操作,避免多任务同时操作导致的判断逻辑错误。
内容的提问来源于stack exchange,提问作者Arda Senoglu
相关产品推荐
相关产品推荐

