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

基于Scrapy的Python-Telegram-Bot如何规避Telegram Flood 429错误?

问题

我监控目标网站的价格下跌情况,当商品价格符合设定的降价规则时,相关信息会被记录到notificate表中。随后通过Scrapy的pipelines.py文件中的代码发送Telegram通知。

有时目标网站会有大量商品降价,notificate表中可能出现200条待发送通知,此时发送会触发Telegram的Flood防护机制,遭遇429错误。

已尝试的方案:

  • Rotating multiple tokens(轮换多令牌)
  • add sleep time(添加等待时间)

但以上方案未解决问题,仍受限于Flood防护机制。

需求:无论notificate表中有多少通知,将其加入队列,并以不超过每秒20条的速率发送。请问该如何实现?

我的pipeline.py代码如下:

def sendnotifications(self, token):
    cursor = self.cnx.cursor()
    req = requests
    cursor.execute("SELECT * FROM notificate WHERE token= '"+token+"'")
    notifications = cursor.fetchall()
    for notification in notifications:
        print(notification)
        productid = notification[1]
        url = notification[3]
        name = notification[2]
        old = notification[4]
        new = notification[5]
        price_difference = old - new
        percentage = price_difference / old
        percentage_str = str("%.2f" % (percentage * 100))
        

        message = "<b>" + name + "</b>" + "\n\n" + \
            str(old) + " TL >>>> " + \
            str(new) + f" TL - {percentage_str}%" + "\n\n" + \
            url + "\n\n" + \
            

        if str(old) == "1.00" or str(old) == "2.00":
            message = "<b>" + name + "</b>" + "\n\n" + \
                "<b>" + str(new) + " TL\n\n" + "</b>" + \
                url + "\n\n" + \
                
            

        token_list = [
            "xxxxxxxxxxxxxxxxxxxxxxx",
            "yyyyyyyyyyyyyyyyyyyyyyyy",
            "zzzzzzzzzzzzzzzzzzzzzzzzzzz",
            "aaaaaaaaaaaaaaaaaaaaaaaaaaaa",
            "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb",
            "ccccccccccccccccccccccccccccccccc",
       

        ]

        TOKEN = token_list[random.randint(0, len(token_list)-1)]


        
        chat_id = "-100xxxxxxxxxxxxx"
        bot = telegram.Bot(token=TOKEN)
        # tel_url = bot.sendMessage(chat_id = chat_id, text = message, parse_mode=ParseMode.HTML)
        
        try:
            bot.sendMessage(chat_id = chat_id, text = message, parse_mode=ParseMode.HTML)
            sleep(0.05)

        except Exception:

            return False
    cursor.close()
    return True
解决方案

要实现严格控制每秒不超过20条的发送速率,同时应对Telegram的Flood防护,需要从队列管理、速率精准控制、错误重试三个方面优化:

1. 用数据库实现持久化队列(替代一次性读取所有数据)

现有代码一次性取出所有通知,数据量大时极易触发限制,改成标记式队列:

  • 给notificate表新增字段:status(可选值:pending/sending/sent/failed)、retry_count(默认0)、next_retry_time(默认NULL)
  • 每次只取出status='pending'且next_retry_time <= NOW()的N条数据(比如每次取20条)
  • 发送前将status设为sending,避免重复读取

2. 用令牌桶算法精准控制发送速率

固定sleep无法适配网络请求耗时,容易出现速率超标或不足的情况,用token-bucket库实现严格速率控制:

from token_bucket import TokenBucket
import time
import telegram

# 初始化令牌桶:每秒产生20个令牌,桶容量20
bucket = TokenBucket(20, 20)

def send_with_rate_limit(bot, chat_id, message):
    # 必须获取令牌才能发送,无令牌则等待
    bucket.consume(1)
    try:
        bot.sendMessage(chat_id=chat_id, text=message, parse_mode=telegram.ParseMode.HTML)
        return True
    except telegram.error.RetryAfter as e:
        # 按Telegram返回的时间等待后重试
        time.sleep(e.retry_after)
        return send_with_rate_limit(bot, chat_id, message)
    except Exception:
        # 标记为失败,稍后重试
        return False

3. 优化令牌轮换与错误处理

  • 提前初始化所有令牌的Bot实例,避免重复创建开销,循环轮换使用:
token_list = [
    "xxxxxxxxxxxxxxxxxxxxxxx",
    "yyyyyyyyyyyyyyyyyyyyyyyy",
    "zzzzzzzzzzzzzzzzzzzzzzzzzzz",
    "aaaaaaaaaaaaaaaaaaaaaaaaaaaa",
    "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb",
    "ccccccccccccccccccccccccccccccccc",
]
bot_list = [telegram.Bot(token=token) for token in token_list]
current_bot_idx = 0

def get_next_bot():
    global current_bot_idx
    bot = bot_list[current_bot_idx]
    current_bot_idx = (current_bot_idx + 1) % len(bot_list)
    return bot
  • 遇到429错误时,不要直接终止所有发送,根据Telegram返回的RetryAfter值等待后重试,或标记该条通知为failed并设置下次重试时间(比如5分钟后)

4. 独立发送进程(解耦Scrapy Pipeline)

Scrapy Pipeline同步执行发送任务会阻塞爬虫,建议解耦:

  • Scrapy Pipeline只负责将符合条件的通知写入notificate表(status='pending')
  • 单独写一个Python脚本,定时(比如每10秒)读取待发送通知,用上述速率控制逻辑发送

完整优化后的核心代码示例

import time
from token_bucket import TokenBucket
import telegram
from telegram import ParseMode
import mysql.connector

class NotificationSender:
    def __init__(self):
        self.cnx = mysql.connector.connect(
            host='your_host',
            user='your_user',
            password='your_pass',
            database='your_db'
        )
        self.token_list = [
            "xxxxxxxxxxxxxxxxxxxxxxx",
            "yyyyyyyyyyyyyyyyyyyyyyyy",
            "zzzzzzzzzzzzzzzzzzzzzzzzzzz",
            "aaaaaaaaaaaaaaaaaaaaaaaaaaaa",
            "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb",
            "ccccccccccccccccccccccccccccccccc",
        ]
        self.bot_list = [telegram.Bot(token=token) for token in self.token_list]
        self.current_bot_idx = 0
        # 令牌桶:每秒20条,桶容量20
        self.bucket = TokenBucket(20, 20)

    def get_next_bot(self):
        bot = self.bot_list[self.current_bot_idx]
        self.current_bot_idx = (self.current_bot_idx + 1) % len(self.bot_list)
        return bot

    def format_message(self, notification):
        productid = notification[1]
        url = notification[3]
        name = notification[2]
        old = notification[4]
        new = notification[5]
        price_difference = old - new
        percentage = price_difference / old
        percentage_str = str("%.2f" % (percentage * 100))

        if str(old) in ["1.00", "2.00"]:
            return f"<b>{name}</b>\n\n<b>{new} TL\n\n</b>{url}\n\n"
        else:
            return f"<b>{name}</b>\n\n{old} TL >>>> {new} TL - {percentage_str}%\n\n{url}\n\n"

    def send_single_notification(self, notification_id, message):
        bot = self.get_next_bot()
        chat_id = "-100xxxxxxxxxxxxx"
        self.bucket.consume(1)
        try:
            bot.sendMessage(chat_id=chat_id, text=message, parse_mode=ParseMode.HTML)
            # 标记为已发送
            cursor = self.cnx.cursor()
            cursor.execute("UPDATE notificate SET status='sent' WHERE id=%s", (notification_id,))
            self.cnx.commit()
            cursor.close()
            return True
        except telegram.error.RetryAfter as e:
            # 按Telegram要求等待后重试
            time.sleep(e.retry_after)
            return self.send_single_notification(notification_id, message)
        except Exception as e:
            # 标记为失败,增加重试次数,设置下次重试时间
            cursor = self.cnx.cursor()
            cursor.execute("""
                UPDATE notificate 
                SET status='failed', retry_count=retry_count+1, next_retry_time=DATE_ADD(NOW(), INTERVAL 5 MINUTE) 
                WHERE id=%s
            """, (notification_id,))
            self.cnx.commit()
            cursor.close()
            return False

    def process_queue(self):
        cursor = self.cnx.cursor()
        # 每次取20条待发送或重试的通知
        cursor.execute("""
            SELECT id, productid, name, url, old_price, new_price 
            FROM notificate 
            WHERE (status='pending' OR (status='failed' AND next_retry_time <= NOW()))
            LIMIT 20
        """)
        notifications = cursor.fetchall()
        cursor.close()

        for notify in notifications:
            notify_id = notify[0]
            message = self.format_message(notify)
            self.send_single_notification(notify_id, message)

if __name__ == "__main__":
    sender = NotificationSender()
    # 循环处理队列,可配合定时任务使用
    while True:
        sender.process_queue()
        time.sleep(1)

关键注意点

  • 必须处理telegram.error.RetryAfter异常,这是Telegram明确要求的等待时间,无视会导致临时封禁
  • 分批次读取通知,避免一次性加载大量数据占用内存
  • 令牌桶算法比固定sleep更精准,能适配网络请求耗时,确保每秒发送量不超限
  • 解耦Scrapy与发送逻辑,避免爬虫被发送任务阻塞

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 20:45:30