如何基于aiogram实现Telegram Bot订单模块的异步定时关闭方案?
基于aiogram的批量定时订单关闭优化方案
针对大量用户订单的定时关闭需求,直接为每个订单创建独立的asyncio.sleep任务会导致资源占用过高,以下是基于异步数据库+定时扫描+任务调度的优化实现方案:
核心思路
- 持久化存储订单:将订单信息(含关闭时间、状态)存入异步数据库,确保机器人重启后不丢失任务。
- 定时批量扫描:启动一个异步后台任务,周期性扫描数据库,筛选出即将到期的订单。
- 异步任务延迟执行:对每个到期订单,计算延迟时间后提交异步任务执行关闭逻辑,避免阻塞主事件循环。
- 时区统一处理:所有时间统一使用UTC存储,避免时区差异导致的关闭时间误差。
具体实现步骤
1. 定义异步数据库模型
使用SQLAlchemy异步版实现订单表,存储关键信息:
from sqlalchemy import Column, Integer, DateTime, Boolean from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine from sqlalchemy.orm import declarative_base Base = declarative_base() class Order(Base): __tablename__ = "orders" id = Column(Integer, primary_key=True, autoincrement=True) user_id = Column(Integer, nullable=False) # 关联用户ID close_time = Column(DateTime, nullable=False) # UTC时间 is_closed = Column(Boolean, default=False) # 订单状态 # 初始化异步数据库引擎(以PostgreSQL为例) engine = create_async_engine("postgresql+asyncpg://username:password@localhost/db_name")
2. 实现定时扫描与订单关闭逻辑
创建后台扫描任务,批量处理即将到期的订单:
import asyncio from datetime import datetime, timedelta from sqlalchemy import select from aiogram import Bot # 全局Bot实例(需在main函数中初始化) bot: Bot = None async def scan_pending_orders(session: AsyncSession): while True: now = datetime.utcnow() # 扫描接下来1分钟内要关闭的未完成订单 time_window_end = now + timedelta(minutes=1) stmt = select(Order).where( Order.is_closed == False, Order.close_time >= now, Order.close_time <= time_window_end ) result = await session.execute(stmt) pending_orders = result.scalars().all() for order in pending_orders: # 计算延迟执行的秒数 delay = (order.close_time - now).total_seconds() # 提交异步关闭任务 asyncio.create_task(close_order(session, order.id, delay)) # 每分钟扫描一次,可根据订单量调整间隔 await asyncio.sleep(60) async def close_order(session: AsyncSession, order_id: int, delay: float): await asyncio.sleep(delay) # 二次校验订单状态,避免重复处理 stmt = select(Order).where(Order.id == order_id) result = await session.execute(stmt) order = result.scalar_one_or_none() if not order or order.is_closed: return # 标记订单为已关闭 order.is_closed = True await session.commit() # 通知用户订单关闭(示例) try: await bot.send_message(order.user_id, "您的订单已自动关闭") except Exception as e: # 处理消息发送失败的情况,比如用户已拉黑机器人 print(f"通知用户 {order.user_id} 失败: {str(e)}")
3. 集成到aiogram启动流程
在机器人启动时初始化数据库并启动扫描任务:
from aiogram import Dispatcher from aiogram.types import Message from aiogram.filters import Command async def on_startup(dp: Dispatcher): # 创建数据库表(首次运行时) async with engine.begin() as conn: await conn.run_sync(Base.metadata.create_all) # 启动订单扫描任务 dp.loop.create_task(scan_pending_orders(AsyncSession(bind=engine))) async def handle_create_order(message: Message): # 假设用户输入格式为"2022-08-17 19:00",需先解析成本地时间再转UTC # 这里需根据实际交互逻辑处理用户输入的日期和时间 input_time_str = message.text.strip() try: # 替换为用户所在时区,比如Asia/Shanghai local_time = datetime.strptime(input_time_str, "%Y-%m-%d %H:%M").replace(tzinfo=pytz.timezone("Asia/Shanghai")) utc_close_time = local_time.astimezone(pytz.utc).replace(tzinfo=None) except ValueError: await message.reply("请输入正确的时间格式:YYYY-MM-DD HH:MM") return # 存入数据库 async with AsyncSession(bind=engine) as session: new_order = Order(user_id=message.from_user.id, close_time=utc_close_time) session.add(new_order) await session.commit() await message.reply(f"订单已创建,将在 {input_time_str} 自动关闭") async def main(): global bot bot = Bot(token="YOUR_TELEGRAM_BOT_TOKEN") dp = Dispatcher() # 注册启动钩子和消息处理器 dp.startup.register(on_startup) dp.message.register(handle_create_order, Command("create_order")) await dp.start_polling(bot) if __name__ == "__main__": import pytz asyncio.run(main())
关键优化点
- 资源高效利用:通过批量扫描替代单个订单的独立睡眠任务,减少并发任务数量,降低内存占用。
- 任务可靠性:数据库持久化确保机器人重启后未完成的订单不会丢失,扫描任务会自动续接。
- 并发安全:二次校验订单状态,配合数据库事务,避免重复关闭订单。
- 容错处理:在消息通知环节添加异常捕获,避免单个任务失败影响全局。
内容的提问来源于stack exchange,提问作者loldude
相关产品推荐
相关产品推荐

