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

Telegram机器人视频编辑队列JSON文件同步异常问题求助

Telegram机器人JSON队列同步问题解决方案

问题概述

开发的Telegram机器人负责接收用户视频,完成下载、编辑后返回成品。基于JSON文件实现任务队列,记录等待/正在编辑的用户ID,编辑完成后用pop()移除队首元素。但在并发场景下(比如用户完成出队同时新用户入队),出现数据覆盖、队列逻辑混乱的问题,即使写入前读取文件也无法解决。

核心问题分析

  1. 无原子性保障:异步任务并发读写JSON文件时,读和写操作被拆分,中间可能被其他任务插入,导致数据覆盖。比如任务A读取队列后,任务B也读取并修改,A再写入就会覆盖B的更改。
  2. 本地缓存过期:循环等待队列时依赖本地的queue变量,没有实时读取文件,导致判断的不是最新队列状态。
  3. 重复冗余的读写逻辑:代码中多次重复读写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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 03:59:59