Telegram机器人视频编辑队列JSON文件同步异常问题求助
Telegram机器人JSON队列同步问题解决方案
问题概述
开发的Telegram机器人负责接收用户视频,完成下载、编辑后返回成品。基于JSON文件实现任务队列,记录等待/正在编辑的用户ID,编辑完成后用pop()移除队首元素。但在并发场景下(比如用户完成出队同时新用户入队),出现数据覆盖、队列逻辑混乱的问题,即使写入前读取文件也无法解决。
核心问题分析
- 无原子性保障:异步任务并发读写JSON文件时,读和写操作被拆分,中间可能被其他任务插入,导致数据覆盖。比如任务A读取队列后,任务B也读取并修改,A再写入就会覆盖B的更改。
- 本地缓存过期:循环等待队列时依赖本地的
queue变量,没有实时读取文件,导致判断的不是最新队列状态。 - 重复冗余的读写逻辑:代码中多次重复读写JSON文件,不仅冗余,还增加了并发冲突的概率。
优化方案
1. 引入文件锁实现原子操作
使用filelock库的FileLock,确保每个队列操作(读、写、修改)都是原子性的,同一时间只有一个任务能操作JSON文件。
2. 封装队列操作函数
把队列的增、改、查、删封装成独立函数,统一处理锁和文件读写,避免重复代码,降低出错概率。
3. 修正队列等待逻辑
循环等待时实时读取最新的队列数据,避免依赖本地缓存的过期数据。
4. 简化用户ID标识规则
统一用户ID的格式,比如用{user_id}-{count}的格式标识重复提交的任务,避免复杂的字符串截取判断。
修改后的完整代码
import asyncio import json import os import math from filelock import FileLock # 全局锁路径,和队列文件同目录 queue_lock_path = f"{database_folder}/queue.json.lock" def get_queue(): """原子性读取队列""" queue_path = f"{database_folder}/queue.json" with FileLock(queue_lock_path): if not os.path.isfile(queue_path): with open(queue_path, 'w') as f: json.dump({"queue": []}, f) with open(queue_path, 'r') as f: return json.load(f)["queue"] def update_queue(new_queue): """原子性更新队列""" queue_path = f"{database_folder}/queue.json" with FileLock(queue_lock_path): with open(queue_path, 'w') as f: json.dump({"queue": new_queue}, f) async def edit_video(user_id, profile, message): user_id_str = str(user_id) queue = get_queue() # 处理用户重复入队的情况 # 统计当前用户的任务数量 user_task_count = sum(1 for item in queue if item.split('-')[0] == user_id_str) if user_task_count > 0: new_task_id = f"{user_id_str}-{user_task_count + 1}" queue.append(new_task_id) else: queue.append(user_id_str) update_queue(queue) # 获取用户的队列位置 # 找到所有属于该用户的任务,取第一个的位置 user_positions = [idx + 1 for idx, item in enumerate(queue) if item.split('-')[0] == user_id_str] await message.reply(f"已为您加入队列!➡️ 位置:{user_positions[0]}/{len(queue)} 预计等待时间:{math.trunc(len(queue)*0.7)} 分钟") # 等待队列轮到自己 while True: queue = get_queue() if not queue: await asyncio.sleep(5) continue first_item = queue[0] if first_item.endswith("-editing"): await asyncio.sleep(5) continue first_user_id = first_item.split('-')[0] if first_user_id == user_id_str: # 标记为正在编辑 queue[0] = f"{first_item}-editing" update_queue(queue) break await asyncio.sleep(5) try: random_id = str(id_generator()) video_file = await message.reply_to_message.video.get_file() checkfolders(message.chat.id, profile) temp_download_path = f"{download_folder}/downloaded_temp/{random_id}-{user_id_str}.mp4" await video_file.download(temp_download_path) with open(f"{download_folder}/downloaded_temp/{random_id}-{user_id_str}.txt", "w") as f: f.write(message.text) await bot.send_message(chat_id=user_id, text=f"正在为您制作{profile}风格的视频...") output = await run_external_script(user_id, profile, random_id, temp_download_path) if output: with open(output[-1], "rb") as video_out: await bot.send_video(chat_id=user_id, video=video_out) else: print("外部脚本未返回输出文件路径") # 完成后弹出队首元素 queue = get_queue() if queue: queue.pop(0) update_queue(queue) except Exception as e: print(f"处理视频时出错:{str(e)}") # 出错后也要移除队列中的任务,避免阻塞队列 queue = get_queue() if queue and queue[0].split('-')[0] == user_id_str: queue.pop(0) update_queue(queue)
额外说明
- 需要先安装
filelock库:pip install filelock - 锁文件和队列文件放在同一目录,确保锁的有效性
- 出错后自动移除队列中的任务,避免队列被异常任务阻塞
- 统一用
split('-')[0]提取用户ID,避免字符串截取长度出错的问题
内容的提问来源于stack exchange,提问作者Pady
相关产品推荐
相关产品推荐

