从零构建Agent定时调度系统:Cron表达式、持久化与高可用实践
1. 项目概述:为什么你的Agent需要一个“日程表”?
最近在折腾Agent开发的朋友,估计都遇到过这么个场景:你精心调教的Agent,能帮你查天气、写周报、分析数据,但每次都得你手动去“戳”它一下。比如,你想让它每天早上9点自动整理前一天的销售数据并发到群里,或者每周五下午5点提醒团队提交周报。这时候,你就会发现,一个只会被动响应的Agent,就像一个没有日程表的助理,能力再强也显得不够“聪明”。
这就是我们今天要聊的核心:为Agent赋予定时调度的能力。简单说,就是教会你的Agent看“闹钟”,让它能自主、准时地执行未来某个时间点的任务。这不仅仅是加个setTimeout那么简单,它涉及到任务的定义、存储、触发、执行以及异常处理等一系列复杂问题,是Agent从“工具”迈向“自动化助手”的关键一步。
从网络上的讨论热度也能看出,无论是“Cron表达式”的频繁出现,还是对“持久化”、“通知机制”的关注,都指向了开发者们在构建实用Agent时遇到的共同痛点。一个健壮的定时调度系统,能让你的Agent项目真正落地,从玩具变成生产力工具。接下来,我就结合自己的踩坑经验,拆解如何从零搭建一个可靠、易扩展的Agent定时调度系统。
2. 核心需求与架构设计解析
在动手写代码之前,我们必须想清楚这个调度系统到底要解决哪些问题。拍脑袋设计,后面大概率要返工。
2.1 调度系统的四大核心需求
根据我过去在多个Agent项目中集成调度功能的经验,一个合格的调度系统至少要满足以下四点:
精准与灵活的时间控制:这是基础。我们需要支持像“每天凌晨1点”、“每周一上午10点”、“每隔25分钟”这样的复杂时间规则。Cron表达式几乎是行业标准,它足够强大和通用。前端展示和配置Cron表达式是个小难点,但社区有成熟的组件库可以解决。
任务的持久化与状态管理:Agent服务可能会重启,内存中的任务队列不能丢。我们必须把任务定义(做什么、何时做)持久化到数据库或文件中。更重要的是,任务本身有生命周期(等待中、执行中、成功、失败、已取消),我们需要可靠地记录和更新这些状态,防止任务重复执行或丢失。
可靠的通知与回调机制:任务执行完了,成功或失败,总得有个说法。系统需要能将结果通知给相关方。这不仅仅是发个日志那么简单,可能需要回调某个HTTP接口、发送消息到钉钉/飞书群,或者触发另一个Agent工作流。这是实现多Agent协作和复杂业务流程的关键。
优雅的异常处理与容错:任务执行时网络超时、依赖服务挂了、Agent本身逻辑出Bug……这些情况太常见了。调度系统不能因为一个任务崩溃就导致整个调度器瘫痪。我们需要有失败重试机制、超时控制,并且能记录详细的错误日志便于排查。
2.2 主流技术方案选型与对比
明确了需求,我们来看看有哪些轮子可以用,以及为什么我最终推荐自研一个轻量级核心。
方案一:直接使用成熟的调度框架比如Java界的Quartz,或者Python的APScheduler。它们功能非常强大,开箱即用,持久化、集群、故障转移都支持。
- 优点:省心,稳定,适合大型复杂系统。
- 缺点:重,与Agent框架(如LangChain、Semantic Kernel等)的集成需要额外封装;对于追求轻量、高定制化的Agent项目来说,可能有些功能用不上,反而增加了复杂度。
方案二:基于云服务或K8s CronJob如果你部署在云上,可以直接用云函数(如AWS Lambda)的定时触发器,或者Kubernetes的CronJob。
- 优点:无需管理调度器本身,利用云原生能力,伸缩性好。
- 缺点:将调度逻辑与业务逻辑(Agent)分离了,任务状态跟踪、跨任务数据传递变得困难;也受限于特定云厂商。
方案三:自研轻量级调度核心这是我个人在中小型Agent项目中更倾向的方案。核心很简单:一个解析Cron表达式的库 + 一个持久化存储 + 一个常驻后台线程/进程去扫描和执行。
- 优点:极度轻量,与Agent业务逻辑无缝集成,定制自由度极高,可以完美适配你的Agent框架和通知机制。
- 缺点:需要自己实现可靠性保障(如分布式锁、故障恢复),适合对系统有较强掌控力的开发者。
对于大多数Agent开发进阶者,我建议从方案三开始。它能让你透彻理解调度系统的每一个环节,而且现代语言(如Python的schedule库结合croniter,或Node.js的node-cron)已经让这件事变得非常简单。下面,我们就以Python为例,搭建这个核心。
3. 核心模块设计与实现细节
我们来把调度系统拆解成几个核心模块,一个个实现。
3.1 任务定义与数据模型设计
首先,我们需要一个数据结构来完整描述一个定时任务。这将是存储在数据库里的核心。
from pydantic import BaseModel, Field from datetime import datetime from enum import Enum from typing import Any, Dict, Optional class TaskStatus(str, Enum): PENDING = "pending" # 等待执行 RUNNING = "running" # 执行中 SUCCESS = "success" # 成功 FAILED = "failed" # 失败 CANCELLED = "cancelled" # 已取消 class ScheduledTask(BaseModel): """定时任务数据模型""" task_id: str = Field(..., description="任务唯一ID") name: str = Field(..., description="任务名称,如'每日销售报告'") # 核心:Cron表达式,定义执行时间 cron_expression: str = Field(..., description="Cron表达式,如 '0 9 * * *' 表示每天9点") # 任务负载:这里定义Agent要执行的动作 agent_action: str = Field(..., description="Agent执行的动作标识,如 'generate_daily_report'") action_payload: Dict[str, Any] = Field(default_factory=dict, description="传递给Action的参数") # 状态与元数据 status: TaskStatus = Field(default=TaskStatus.PENDING, description="当前状态") last_run_time: Optional[datetime] = Field(default=None, description="上次执行时间") next_run_time: Optional[datetime] = Field(default=None, description="下次预计执行时间") created_at: datetime = Field(default_factory=datetime.now) updated_at: datetime = Field(default_factory=datetime.now) # 失败重试配置 retry_count: int = Field(default=0, description="已重试次数") max_retries: int = Field(default=3, description="最大重试次数") # 通知配置(可选) webhook_url: Optional[str] = Field(default=None, description="任务完成后的回调URL") notify_on_failure: bool = Field(default=True, description="是否在失败时通知")设计要点解析:
agent_action与action_payload:这是连接调度系统与Agent业务逻辑的桥梁。agent_action可以是你Agent内部注册的一个函数名或技能(Skill)名,payload则是调用时需要的参数。这种设计实现了调度与业务的解耦。next_run_time的预计算:为了高效扫描即将执行的任务,我们应在任务创建或每次执行后,立即根据Cron表达式计算出下一次运行时间并存储。这样调度器只需查询next_run_time <= now()的任务即可,避免每次都对所有任务的Cron表达式进行解析计算。- 使用Pydantic:利用Pydantic进行数据验证和序列化,能省去很多手动检查的代码,尤其在与API或数据库交互时非常方便。
3.2 持久化存储层实现
任务数据必须持久化。这里我选择SQLite作为起步,因为它无需额外服务,简单可靠。后期可以轻松迁移到PostgreSQL或MySQL。
import sqlite3 from contextlib import contextmanager from typing import List, Optional, Generator import json class TaskStore: """任务存储层,负责任务的CRUD和状态持久化""" def __init__(self, db_path: str = "agent_scheduler.db"): self.db_path = db_path self._init_db() def _init_db(self): """初始化数据库表""" with self._get_connection() as conn: conn.execute(""" CREATE TABLE IF NOT EXISTS scheduled_tasks ( task_id TEXT PRIMARY KEY, name TEXT NOT NULL, cron_expression TEXT NOT NULL, agent_action TEXT NOT NULL, action_payload TEXT NOT NULL, -- 存储为JSON字符串 status TEXT NOT NULL, last_run_time TIMESTAMP, next_run_time TIMESTAMP, created_at TIMESTAMP NOT NULL, updated_at TIMESTAMP NOT NULL, retry_count INTEGER DEFAULT 0, max_retries INTEGER DEFAULT 3, webhook_url TEXT, notify_on_failure BOOLEAN DEFAULT 1 ) """) # 为 next_run_time 创建索引,加速扫描查询 conn.execute("CREATE INDEX IF NOT EXISTS idx_next_run_time ON scheduled_tasks(next_run_time)") conn.execute("CREATE INDEX IF NOT EXISTS idx_status ON scheduled_tasks(status)") @contextmanager def _get_connection(self) -> Generator[sqlite3.Connection, None, None]: """获取数据库连接的上下文管理器""" conn = sqlite3.connect(self.db_path, detect_types=sqlite3.PARSE_DECLTYPES) conn.row_factory = sqlite3.Row # 使返回结果为字典式行对象 try: yield conn conn.commit() except Exception: conn.rollback() raise finally: conn.close() def save_task(self, task: ScheduledTask) -> str: """创建或更新任务""" with self._get_connection() as conn: payload_json = json.dumps(task.action_payload) conn.execute(""" INSERT OR REPLACE INTO scheduled_tasks (task_id, name, cron_expression, agent_action, action_payload, status, last_run_time, next_run_time, created_at, updated_at, retry_count, max_retries, webhook_url, notify_on_failure) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( task.task_id, task.name, task.cron_expression, task.agent_action, payload_json, task.status, task.last_run_time, task.next_run_time, task.created_at, task.updated_at, task.retry_count, task.max_retries, task.webhook_url, task.notify_on_failure )) return task.task_id def get_due_tasks(self) -> List[ScheduledTask]: """获取所有已到执行时间的任务(next_run_time <= now())""" from datetime import datetime now = datetime.now() tasks = [] with self._get_connection() as conn: cursor = conn.execute( "SELECT * FROM scheduled_tasks WHERE status = ? AND next_run_time <= ?", (TaskStatus.PENDING.value, now) ) for row in cursor: task_dict = dict(row) task_dict['action_payload'] = json.loads(task_dict['action_payload']) # 将数据库中的字符串状态转换回枚举 task_dict['status'] = TaskStatus(task_dict['status']) tasks.append(ScheduledTask(**task_dict)) return tasks def update_task_status(self, task_id: str, status: TaskStatus, last_run_time: Optional[datetime] = None, next_run_time: Optional[datetime] = None, retry_count: Optional[int] = None): """更新任务状态及相关时间""" with self._get_connection() as conn: update_fields = ["status = ?", "updated_at = ?"] params = [status.value, datetime.now()] if last_run_time is not None: update_fields.append("last_run_time = ?") params.append(last_run_time) if next_run_time is not None: update_fields.append("next_run_time = ?") params.append(next_run_time) if retry_count is not None: update_fields.append("retry_count = ?") params.append(retry_count) params.append(task_id) # WHERE 条件参数 update_sql = f"UPDATE scheduled_tasks SET {', '.join(update_fields)} WHERE task_id = ?" conn.execute(update_sql, params)实操心得:
- 索引是关键:务必为
next_run_time和status字段创建索引。当你有成千上万个任务时,没有索引的扫描查询会成为性能瓶颈。 - JSON序列化:将
action_payload这类动态结构存储为JSON字符串,比拆分成多个关系型字段更灵活,更适合Agent任务参数多变的特点。 - 连接管理:使用上下文管理器(
@contextmanager)来管理数据库连接,可以确保连接被正确关闭,即使在发生异常时也能回滚事务,避免数据不一致。
3.3 Cron表达式解析与下次执行时间计算
这是调度器的“大脑”。我们需要一个库来解析Cron表达式并计算下一次触发时间。Python中croniter库是绝佳选择。
from croniter import croniter from datetime import datetime, timedelta class CronScheduler: """处理Cron表达式解析与时间计算""" @staticmethod def get_next_run_time(cron_expression: str, base_time: datetime = None) -> datetime: """根据Cron表达式和基准时间,计算下一次运行时间""" if base_time is None: base_time = datetime.now() try: cron = croniter(cron_expression, base_time) return cron.get_next(datetime) # 返回下一个时间点 except Exception as e: # 这里可以记录日志,并抛出自定义异常 raise ValueError(f"无效的Cron表达式 '{cron_expression}': {e}") @staticmethod def validate_cron_expression(cron_expression: str) -> bool: """验证Cron表达式是否有效""" try: croniter(cron_expression) return True except: return False @staticmethod def get_upcoming_schedule(cron_expression: str, count: int = 5) -> List[datetime]: """获取接下来几次的执行时间,用于调试或展示给用户""" schedule = [] base_time = datetime.now() cron = croniter(cron_expression, base_time) for _ in range(count): schedule.append(cron.get_next(datetime)) return schedule注意事项:
- 时区问题:这是定时任务最常见的坑!
croniter默认使用本地时间。如果你的服务部署在UTC时间的服务器上,而你的Cron表达式是针对北京时间(UTC+8)的,就会产生8小时的偏差。最佳实践是:在存储和计算时,全部使用UTC时间。在创建任务时,根据用户所在时区,将用户输入的“本地时间”Cron表达式,转换为对应的UTC时间Cron表达式再存储。或者,在croniter初始化时指定一个带时区的datetime对象作为基准。 - 表达式校验:一定要在任务创建或更新时校验Cron表达式的有效性,避免无效表达式导致调度器出错。
3.4 调度器核心引擎实现
现在,我们把存储、计算和业务逻辑串联起来,构建调度器的主循环。
import time import threading import logging from concurrent.futures import ThreadPoolExecutor logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s') logger = logging.getLogger(__name__) class AgentScheduler: """Agent定时调度器核心引擎""" def __init__(self, task_store: TaskStore, agent_executor, scan_interval_seconds: int = 30): """ Args: task_store: 任务存储实例 agent_executor: 执行Agent动作的调用器,需实现 execute_action(action, payload) 方法 scan_interval_seconds: 扫描数据库的间隔(秒) """ self.task_store = task_store self.agent_executor = agent_executor self.scan_interval = scan_interval_seconds self._scheduler_thread = None self._stop_event = threading.Event() # 使用线程池执行任务,避免阻塞调度扫描 self._executor = ThreadPoolExecutor(max_workers=5, thread_name_prefix="AgentTaskWorker") self._cron_util = CronScheduler() def start(self): """启动调度器""" if self._scheduler_thread and self._scheduler_thread.is_alive(): logger.warning("调度器已在运行中") return self._stop_event.clear() self._scheduler_thread = threading.Thread(target=self._run_scheduler_loop, name="SchedulerMainLoop") self._scheduler_thread.daemon = True # 设置为守护线程,主程序退出时自动结束 self._scheduler_thread.start() logger.info("Agent定时调度器已启动") def stop(self): """停止调度器""" logger.info("正在停止调度器...") self._stop_event.set() if self._scheduler_thread: self._scheduler_thread.join(timeout=10) # 等待线程结束 self._executor.shutdown(wait=True) logger.info("调度器已停止") def _run_scheduler_loop(self): """调度器主循环""" logger.info("调度器主循环开始运行") while not self._stop_event.is_set(): try: self._scan_and_execute_tasks() except Exception as e: logger.error(f"调度器主循环发生未预期错误: {e}", exc_info=True) # 等待指定间隔,但可被停止事件中断 self._stop_event.wait(self.scan_interval) logger.info("调度器主循环结束") def _scan_and_execute_tasks(self): """扫描并执行到期任务的核心方法""" # 1. 从数据库获取所有已到期的任务 due_tasks = self.task_store.get_due_tasks() if not due_tasks: return logger.info(f"扫描到 {len(due_tasks)} 个到期任务") for task in due_tasks: # 2. 立即将任务状态更新为“执行中”,防止被其他进程/线程重复捞取 # 注意:在分布式环境下,这里需要更严谨的分布式锁,例如基于Redis的锁 self.task_store.update_task_status( task_id=task.task_id, status=TaskStatus.RUNNING, last_run_time=datetime.now() # 记录开始执行时间 ) # 3. 将任务提交到线程池异步执行 future = self._executor.submit(self._execute_single_task, task) # 可以添加回调来处理执行结果,这里我们简化处理,在_execute_single_task内部更新状态 def _execute_single_task(self, task: ScheduledTask): """在独立线程中执行单个任务""" task_id = task.task_id logger.info(f"开始执行任务: {task.name} (ID: {task_id})") try: # 1. 调用Agent执行器,执行业务逻辑 # 这里的 agent_executor 需要你根据具体的Agent框架实现 # 例如:result = self.agent_executor.execute_action(task.agent_action, task.action_payload) # 为了演示,我们模拟一个执行过程 result = self._simulate_agent_execution(task) # 2. 计算下一次执行时间 next_run = self._cron_util.get_next_run_time(task.cron_expression) # 3. 更新任务状态为成功,并设置下次执行时间 self.task_store.update_task_status( task_id=task_id, status=TaskStatus.SUCCESS, next_run_time=next_run, retry_count=0 # 成功则重置重试计数 ) logger.info(f"任务执行成功: {task.name} (ID: {task_id})") # 4. 发送成功通知(如果配置了webhook) if task.webhook_url: self._send_notification(task, success=True, result=result) except Exception as e: logger.error(f"任务执行失败: {task.name} (ID: {task_id}), 错误: {e}", exc_info=True) # 处理失败逻辑 self._handle_task_failure(task, e) def _simulate_agent_execution(self, task: ScheduledTask): """模拟Agent执行过程,实际项目中替换为真实的Agent调用""" # 模拟一个耗时操作 time.sleep(2) return {"message": f"模拟执行动作 {task.agent_action} 成功", "data": task.action_payload} def _handle_task_failure(self, task: ScheduledTask, error: Exception): """处理任务执行失败""" current_retry = task.retry_count + 1 if current_retry < task.max_retries: # 还可以重试 logger.info(f"任务 {task.name} 准备第 {current_retry} 次重试") # 可以设置一个退避延迟,比如 2^retry_count 分钟后再试 delay_minutes = 2 ** current_retry next_retry_time = datetime.now() + timedelta(minutes=delay_minutes) # 将任务状态改回PENDING,并设置一个近期的next_run_time用于重试 # 注意:这里修改了next_run_time,会覆盖原有的Cron计划。重试是临时调度。 self.task_store.update_task_status( task_id=task.task_id, status=TaskStatus.PENDING, next_run_time=next_retry_time, retry_count=current_retry ) else: # 重试次数用尽,标记为失败 logger.error(f"任务 {task.name} 重试次数用尽,标记为永久失败") next_run = self._cron_util.get_next_run_time(task.cron_expression) # 仍然计算下一次常规执行时间 self.task_store.update_task_status( task_id=task.task_id, status=TaskStatus.FAILED, next_run_time=next_run, retry_count=current_retry ) # 发送失败通知 if task.webhook_url or task.notify_on_failure: self._send_notification(task, success=False, error=str(error)) def _send_notification(self, task: ScheduledTask, success: bool, result=None, error=None): """发送任务执行结果通知(例如调用Webhook)""" # 这里可以实现HTTP请求到配置的webhook_url # 或者集成消息通知服务(如钉钉机器人、飞书机器人、邮件等) # 示例:使用requests库发送POST请求 import requests payload = { "task_id": task.task_id, "task_name": task.name, "status": "success" if success else "failed", "timestamp": datetime.now().isoformat(), "result": result, "error": error } try: # 注意:在实际发送前,请确认webhook_url有效且安全 # response = requests.post(task.webhook_url, json=payload, timeout=5) # response.raise_for_status() logger.info(f"已发送任务通知: {task.name}, 状态: {payload['status']}") except Exception as e: logger.error(f"发送任务通知失败: {e}")核心逻辑拆解:
- 主循环 (
_run_scheduler_loop):一个简单的while循环,定期扫描数据库。使用threading.Event的wait方法来实现可中断的睡眠,这样在调用stop()时能快速退出。 - 任务状态机:这是保证系统可靠性的关键。一个任务从
PENDING被扫描到,立即置为RUNNING,执行完毕后根据结果转为SUCCESS或FAILED。状态转换必须在持久化层原子性完成,防止并发执行。 - 异步执行:使用
ThreadPoolExecutor将任务执行与扫描解耦。扫描线程不会被耗时的Agent任务阻塞,可以继续发现新的到期任务。线程池大小 (max_workers) 需要根据你Agent任务的IO/CPU密集程度和系统资源来调整。 - 失败重试与退避:
_handle_task_failure实现了简单的指数退避重试。失败后不是立即重试,而是等待一段时间(如2分钟、4分钟、8分钟),避免在服务瞬时故障时产生雪崩效应。
4. 与Agent框架的集成实践
调度器是“骨架”,现在需要注入“灵魂”——让你的Agent真正动起来。这里的关键是agent_executor。
4.1 定义统一的Agent动作执行接口
首先,我们需要一个抽象层,让调度器能以统一的方式调用不同的Agent能力。
from abc import ABC, abstractmethod from typing import Any, Dict class AgentActionExecutor(ABC): """Agent动作执行器抽象接口""" @abstractmethod def execute_action(self, action_name: str, payload: Dict[str, Any]) -> Any: """ 执行指定的Agent动作。 Args: action_name: 动作标识符,如 "send_email", "analyze_data" payload: 动作所需的参数 Returns: 动作执行的结果,可以是任何可序列化的对象 Raises: ActionNotFoundException: 当action_name未注册时 ActionExecutionFailedException: 当动作执行过程中出错时 """ pass4.2 实现基于流行Agent框架的执行器
假设你的Agent是基于LangChain或类似框架构建的,下面是一个集成示例:
import importlib from typing import Callable class SimpleAgentExecutor(AgentActionExecutor): """一个简单的、基于函数注册的Agent执行器""" def __init__(self): self._action_registry = {} # 存储 action_name -> 可调用函数 的映射 def register_action(self, action_name: str, action_func: Callable): """向执行器注册一个动作函数""" if action_name in self._action_registry: logger.warning(f"动作 '{action_name}' 已存在,将被覆盖") self._action_registry[action_name] = action_func logger.info(f"已注册动作: {action_name}") def execute_action(self, action_name: str, payload: Dict[str, Any]) -> Any: if action_name not in self._action_registry: raise ValueError(f"未找到注册的动作: {action_name}") action_func = self._action_registry[action_name] logger.info(f"正在执行动作: {action_name}, 参数: {payload}") # 调用注册的函数,并传入参数 return action_func(**payload) # 示例:注册几个具体的Agent技能(Skill) def generate_daily_report(date: str = None, format: str = "markdown"): """生成每日报告 - 这是一个模拟的Agent技能""" from datetime import datetime, timedelta if not date: date = (datetime.now() - timedelta(days=1)).strftime('%Y-%m-%d') # 这里应该是你真实的报告生成逻辑,例如调用LLM、查询数据库等 report_content = f"# {date} 每日运营报告\n\n- 模拟生成报告内容...\n- 格式要求: {format}" logger.info(f"已生成 {date} 的报告") return {"report_date": date, "content": report_content, "format": format} def send_team_reminder(channel: str, message: str): """发送团队提醒 - 另一个模拟技能""" # 这里可以集成钉钉、飞书、企业微信等消息机器人 logger.info(f"模拟发送提醒到 [{channel}]: {message}") return {"status": "sent", "channel": channel, "message": message} # 初始化执行器并注册技能 agent_executor = SimpleAgentExecutor() agent_executor.register_action("generate_daily_report", generate_daily_report) agent_executor.register_action("send_team_reminder", send_team_reminder) # 现在,调度器可以这样调用: # result = agent_executor.execute_action("generate_daily_report", {"date": "2023-10-27"})集成模式扩展:
- 与LangChain集成:你的
action_func可以是一个LangChain的Chain或Agent的invoke方法。 - 与Semantic Kernel集成:可以注册Semantic Function或Native Function作为可调度动作。
- 与HTTP服务集成:如果你的Agent能力以HTTP API形式暴露,
action_func可以封装一个HTTP客户端调用。
4.3 创建与管理定时任务
最后,我们需要一个管理面,来创建、查看、更新和删除定时任务。这通常以一个HTTP API或命令行工具的形式提供。
from uuid import uuid4 from fastapi import FastAPI, HTTPException, BackgroundTasks # 假设使用FastAPI构建API from pydantic import BaseModel app = FastAPI(title="Agent定时调度系统API") # 依赖注入(在实际应用中,这些应该通过依赖注入框架管理) task_store = TaskStore() scheduler = AgentScheduler(task_store, agent_executor) scheduler.start() # 启动调度器 class CreateTaskRequest(BaseModel): name: str cron_expression: str agent_action: str action_payload: Dict[str, Any] = {} max_retries: int = 3 webhook_url: Optional[str] = None notify_on_failure: bool = True @app.post("/tasks") async def create_task(req: CreateTaskRequest): """创建新的定时任务""" # 1. 验证Cron表达式 if not CronScheduler.validate_cron_expression(req.cron_expression): raise HTTPException(status_code=400, detail="无效的Cron表达式") # 2. 验证Agent动作是否存在 # 这里可以添加检查,确保 req.agent_action 已在 executor 中注册 # 3. 创建任务对象 task_id = str(uuid4()) next_run_time = CronScheduler.get_next_run_time(req.cron_expression) new_task = ScheduledTask( task_id=task_id, name=req.name, cron_expression=req.cron_expression, agent_action=req.agent_action, action_payload=req.action_payload, status=TaskStatus.PENDING, next_run_time=next_run_time, max_retries=req.max_retries, webhook_url=req.webhook_url, notify_on_failure=req.notify_on_failure ) # 4. 保存到数据库 task_store.save_task(new_task) # 5. 立即触发一次调度扫描(可选),让新任务如果立即到期也能被快速执行 # background_tasks.add_task(scheduler._scan_and_execute_tasks) return {"task_id": task_id, "message": "任务创建成功", "next_run_time": next_run_time} @app.get("/tasks/{task_id}") async def get_task(task_id: str): """获取任务详情""" # 实现从数据库查询的逻辑... pass @app.delete("/tasks/{task_id}") async def cancel_task(task_id: str): """取消(删除)一个任务""" # 实现将任务状态更新为 CANCELLED 或直接从数据库删除的逻辑... pass @app.get("/tasks") async def list_tasks(status: Optional[TaskStatus] = None, page: int = 1, size: int = 20): """分页列出所有任务,可按状态过滤""" # 实现数据库分页查询逻辑... pass通过这样一套API,前端界面或脚本就可以方便地管理Agent的定时任务了。
5. 生产环境进阶考量与优化
上面我们实现了一个可用的单机版调度系统。但要用于生产环境,还需要考虑更多。
5.1 分布式部署与高可用
单点故障是定时任务系统的大忌。我们需要让调度器支持多实例部署。
- 核心矛盾:多个调度器实例同时运行,如何避免同一个任务被重复执行?
- 解决方案:分布式锁。在扫描并获取到期任务 (
get_due_tasks) 后,尝试获取该任务的锁,只有拿到锁的实例才能将其状态改为RUNNING并执行。 - 技术选型:
- 数据库乐观锁/悲观锁:在更新任务状态为
RUNNING时,使用UPDATE ... WHERE status = 'PENDING'这样的原子操作,并检查受影响行数。简单但数据库压力大,且实例间时钟需同步。 - Redis分布式锁:更轻量、性能更好的选择。使用
SETNX(SET if Not eXists) 命令或 Redlock 算法。 - ZooKeeper/etcd:适用于更复杂的协调场景,但重量级。
- 数据库乐观锁/悲观锁:在更新任务状态为
# 伪代码:使用Redis分布式锁 import redis import uuid class DistributedTaskStore(TaskStore): def __init__(self, db_path, redis_client): super().__init__(db_path) self.redis = redis_client self.lock_timeout = 30 # 锁超时时间,秒 def acquire_task_lock(self, task_id: str) -> bool: """尝试获取任务锁,返回是否成功""" lock_key = f"task_lock:{task_id}" lock_value = str(uuid.uuid4()) # 唯一标识当前实例 # 设置锁,NX表示仅当key不存在时设置,EX设置过期时间 acquired = self.redis.set(lock_key, lock_value, nx=True, ex=self.lock_timeout) return acquired is not None def release_task_lock(self, task_id: str, lock_value: str): """释放任务锁,使用Lua脚本保证原子性""" lock_key = f"task_lock:{task_id}" lua_script = """ if redis.call("get", KEYS[1]) == ARGV[1] then return redis.call("del", KEYS[1]) else return 0 end """ self.redis.eval(lua_script, 1, lock_key, lock_value)在_scan_and_execute_tasks方法中,在更新任务状态前,先尝试获取锁。
5.2 任务执行的可观测性
出了问题得能快速定位。
- 结构化日志:使用像
structlog或json-logger这样的库,为每条日志记录任务ID、动作名称、执行时间等上下文信息。 - 指标监控:集成Prometheus等监控系统,暴露指标如
scheduler_tasks_total(任务总数)、scheduler_tasks_executed(已执行)、scheduler_tasks_failed(失败)、scheduler_task_duration_seconds(执行耗时直方图)。 - 链路追踪:为每个任务执行生成一个唯一的Trace ID,贯穿整个调用链(调度器 -> Agent执行器 -> 外部服务),便于在分布式系统中追踪问题。
5.3 任务依赖与工作流
有时任务不是独立的。比如“任务A:收集数据”必须在“任务B:生成报告”之前完成。
- 简单依赖:可以在任务
action_payload中传递前序任务的ID或结果标识,由Agent逻辑自行判断依赖是否就绪。或者,在调度器中实现一个简单的状态检查,只有前置任务成功,才将本任务状态置为PENDING。 - 复杂工作流:这就需要引入工作流引擎(如Airflow、Prefect)或状态机了。此时,调度器调度的可能不是一个具体的Agent动作,而是一个工作流的启动事件。这属于更高级的架构,需要根据业务复杂度权衡。
5.4 动态配置与热更新
不希望每次修改任务或增减任务都重启服务。
- 我们的设计已支持:通过API创建、更新、删除任务,调度器主循环每次扫描都会从数据库加载最新状态,天然支持动态变更。
- Agent动作热注册:
SimpleAgentExecutor的register_action方法可以在运行时动态添加新的动作,无需重启。
6. 常见问题排查与实战技巧
在实际部署和运行中,我踩过不少坑,这里总结一下。
6.1 任务被重复执行了
- 可能原因1:调度器实例多跑了一个。检查部署流程,确保没有意外启动多个进程。使用
ps aux | grep your_scheduler查看。 - 可能原因2:任务执行时间过长,超过了状态锁的有效期。任务状态已从
RUNNING超时恢复为PENDING,被另一个调度器实例再次捞取。解决:合理设置锁超时时间,应大于任务最大可能执行时间。或者在任务开始执行时,定期“续租”锁。 - 可能原因3:数据库事务隔离级别问题。在极高并发下,两个事务可能同时读到
PENDING状态。解决:使用SELECT ... FOR UPDATE(行级锁)或在更新状态时使用更严格的条件(如WHERE status = 'PENDING' AND version = ?乐观锁)。
6.2 任务没有按时执行,延迟很大
- 可能原因1:调度器扫描间隔 (
scan_interval_seconds) 设置过长。如果设为60秒,那么任务最多可能延迟60秒才被发现。解决:根据业务对准时性的要求调整,例如设为10秒。注意权衡数据库查询频率。 - 可能原因2:线程池已满,任务在队列中等待。如果
max_workers设置过小,而同时到期的任务很多,会导致任务排队。解决:监控线程池队列长度,适当增加max_workers,或者使用有界队列并设置合理的拒绝策略。 - 可能原因3:系统负载过高,CPU调度延迟。解决:监控系统资源,优化Agent任务本身的性能。
6.3 Cron表达式不生效,时间不对
- 首要怀疑对象:时区。这是新手最容易踩的坑。务必明确你的服务器时区、数据库存储的时区、
croniter计算使用的时区。强烈建议全部使用UTC时间,在最终展示给用户时再转换为本地时间。 - 检查Cron表达式语法:使用在线的Cron表达式验证工具,或
CronScheduler.get_upcoming_schedule()方法打印未来几次执行时间,看是否符合预期。 - 检查系统时间:确保服务器时间准确,可以使用NTP服务同步。
6.4 Agent动作执行失败,但错误信息不明确
- 在
_execute_single_task方法中捕获更广泛的异常,并记录完整的堆栈信息 (exc_info=True)。 - 在Agent动作执行器内部做好日志记录,记录输入参数和关键步骤。
- 实现任务执行历史表,不仅记录成功失败,还记录详细的执行日志、错误信息、开始结束时间,便于后期审计和排查。
6.5 如何优雅地停止和重启调度器
- 停止:我们实现的
stop()方法通过设置_stop_event来中断主循环,并等待线程池关闭。确保在程序退出(如接收SIGTERM信号)时调用此方法。 - 重启:由于任务状态和下次执行时间都已持久化,重启调度器服务是安全的。重启后,调度器会从数据库加载所有
PENDING状态的任务,并根据next_run_time继续执行。注意重启期间可能到期的任务,会在重启后的第一次扫描中被执行,这可能导致微小延迟。
为Agent加上定时调度能力,就像给一位能干的助手配上了日历和闹钟,让它从被动响应变为主动规划。这套系统看似复杂,但拆解开来,无非是“任务定义”、“时间计算”、“状态管理”和“可靠执行”几个核心模块。从简单的单机版开始,逐步迭代加入分布式锁、监控、告警,你就能构建出一个支撑起关键业务的自动化Agent体系。
