基于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
相关产品推荐
相关产品推荐

