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

如何基于aiogram实现Telegram Bot订单模块的异步定时关闭方案?

基于aiogram的批量定时订单关闭优化方案

针对大量用户订单的定时关闭需求,直接为每个订单创建独立的asyncio.sleep任务会导致资源占用过高,以下是基于异步数据库+定时扫描+任务调度的优化实现方案:

核心思路

  1. 持久化存储订单:将订单信息(含关闭时间、状态)存入异步数据库,确保机器人重启后不丢失任务。
  2. 定时批量扫描:启动一个异步后台任务,周期性扫描数据库,筛选出即将到期的订单。
  3. 异步任务延迟执行:对每个到期订单,计算延迟时间后提交异步任务执行关闭逻辑,避免阻塞主事件循环。
  4. 时区统一处理:所有时间统一使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 12:18:55