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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 14:35:43