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

如何恢复基于Pyrogram和MongoDB的Telegram机器人广播进程?

解决Telegram广播进程中断后的恢复问题(Pyrogram + MongoDB)

针对GitHub Workflows每6小时重启导致广播中断的问题,核心思路是将广播进度持久化到MongoDB,重启后从断点继续执行。以下是具体实现方案:

1. MongoDB进度存储设计

创建broadcast_tasks集合,用于存储广播任务的状态和进度,文档结构如下:

{
    "_id": ObjectId("唯一任务ID"),
    "task_id": "自定义唯一标识(如UUID)",
    "message_content": "要发送的广播消息内容",
    "unsent_user_ids": [123456, 789012, ...],  # 待发送的用户ID列表
    "sent_user_ids": [345678, ...],             # 已发送成功的用户ID列表
    "failed_user_ids": [901234, ...],           # 发送失败的用户ID列表
    "status": "in_progress",                    # 任务状态:in_progress/finished/failed
    "created_at": ISODate("2024-05-20T10:00:00Z"),
    "updated_at": ISODate("2024-05-20T12:30:00Z")
}

2. 广播任务的进度管理

初始化广播任务

开始广播前,从MongoDB获取所有用户ID,初始化任务并写入数据库:

import uuid
from datetime import datetime
from pymongo import MongoClient
from pyrogram import Client

# 初始化MongoDB连接
client = MongoClient("mongodb://your_host:your_port/")
db = client["your_db_name"]
broadcast_collection = db["broadcast_tasks"]
users_collection = db["users"]  # 假设存储用户的集合

def init_broadcast_task(message_content):
    # 批量获取所有用户ID
    all_users = users_collection.find({}, {"user_id": 1})
    user_ids = [user["user_id"] for user in all_users]
    
    task_id = str(uuid.uuid4())
    task = {
        "task_id": task_id,
        "message_content": message_content,
        "unsent_user_ids": user_ids,
        "sent_user_ids": [],
        "failed_user_ids": [],
        "status": "in_progress",
        "created_at": datetime.utcnow(),
        "updated_at": datetime.utcnow()
    }
    broadcast_collection.insert_one(task)
    return task_id

分批发送并更新进度

采用批量处理+定期更新的方式,避免频繁写入数据库影响性能,每发送一批用户后更新进度:

def run_broadcast(task_id, batch_size=100):
    task = broadcast_collection.find_one({"task_id": task_id, "status": "in_progress"})
    if not task:
        return
    
    unsent_ids = task["unsent_user_ids"]
    app = Client("your_bot_session")
    
    with app:
        while unsent_ids:
            # 取出当前批次的用户ID
            batch = unsent_ids[:batch_size]
            remaining_ids = unsent_ids[batch_size:]
            
            for user_id in batch:
                try:
                    app.send_message(user_id, task["message_content"])
                    task["sent_user_ids"].append(user_id)
                except Exception as e:
                    print(f"发送给用户{user_id}失败: {str(e)}")
                    task["failed_user_ids"].append(user_id)
            
            # 更新数据库中的进度
            broadcast_collection.update_one(
                {"task_id": task_id},
                {
                    "$set": {
                        "unsent_user_ids": remaining_ids,
                        "sent_user_ids": task["sent_user_ids"],
                        "failed_user_ids": task["failed_user_ids"],
                        "updated_at": datetime.utcnow()
                    }
                }
            )
            unsent_ids = remaining_ids
        
        # 广播完成后更新状态
        broadcast_collection.update_one(
            {"task_id": task_id},
            {"$set": {"status": "finished", "updated_at": datetime.utcnow()}}
        )

3. 重启后的任务恢复

机器人启动时,自动检测未完成的广播任务并继续执行:

def resume_pending_broadcasts():
    # 查询所有状态为in_progress的任务
    pending_tasks = broadcast_collection.find({"status": "in_progress"})
    
    for task in pending_tasks:
        print(f"恢复未完成的广播任务: {task['task_id']}")
        run_broadcast(task["task_id"])

# 在机器人启动时优先执行恢复逻辑
if __name__ == "__main__":
    resume_pending_broadcasts()
    # 启动Pyrogram客户端主逻辑
    app = Client("your_bot_session")
    app.run()

4. 适配GitHub Workflows的注意事项

  • 分批加载用户:如果用户量过大(6万+),避免一次性取出所有用户ID,可通过MongoDB的skip()和limit()分批加载到未发送列表,减少内存占用。
  • 缩小批量更新间隔:将batch_size设置为50-100,确保进度更新更频繁,中断后仅损失当前批次的部分工作量。
  • 单独处理失败用户:可对failed_user_ids列表单独发起重试,或标记为待重试状态,后续统一处理。

内容的提问来源于stack exchange,提问作者Shaheen M

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 05:50:27