ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Agent-harness定时任务调度实战:从手动触发到自动执行

Agent-harness定时任务调度实战:从手动触发到自动执行 你有没有过这种经历Agent 开发完了测试的时候跑得挺顺可真上了线每天早上还是得自己手动触发一次或者半夜爬起来看它到底执行完没有。我最早搭 Agent-harness 框架时就是这样的状态直到我把定时任务调度接进去才真正体会到什么叫“Agent 自己会按时上班”。上周我刚好把一个内部用的“每日行业动态简报 Agent”从手动触发改成了定时调度整个过程踩了不少坑也把 Agent-harness 里跟定时调度相关的部分摸了个遍。这篇教程就围绕 Agent-harness 的定时任务调度来写从为什么需要、怎么选型、关键参数怎么配置到完整跑通一个定时 Agent 任务最后附上我实际的排错记录。如果你正准备让 Agent 从“随叫随到”变成“到点就干”这篇文章应该能帮你省掉不少弯路。1. 为什么 Agent 需要一个“闹钟”定时任务调度的本质与场景1.1 从“你叫它干活”到“它到点自己干”先把概念理清楚。Agent-harness 本质上是一个帮助你定义、编排和运行智能体Agent的开发框架它解决的是 Agent 的“工作流”问题Agent 如何感知输入、如何调用工具、如何拆解任务、如何组织最终输出。但框架本身默认不会替你考虑“什么时间执行”这件事。定时任务调度就是在 Agent-harness 的外面再套一个“触发层”让 Agent 的运行入口按照预设的时间规则被自动调用。你可以把它理解成给 Agent 装了一个闹钟到点了闹钟响Agent 起床干活。以前我需要一个“临时任务”就在命令行里敲一条启动命令等它跑完。这是典型的“人肉触发”。但企业里真正有价值的任务往往是重复性的、周期性的每天早上汇总行业资讯每小时巡检一次服务状态每周一生成一份上周数据复盘。这些任务如果你不交给调度器就永远得靠人记着、人触发一旦某天忘了整个数据链路就断了。1.2 Agent-harness 里定时调度的典型应用场景我梳理了一下在 Agent-harness 环境下定时调度最常用的场景有这么几类周期性信息采集与简报生成比如每天早上 8 点调度器触发一个 Agent让它去抓取指定资讯源的关键内容用另一个子 Agent 做摘要提炼最后生成一份结构化简报推送到钉钉、飞书或者邮件。定时巡检与异常上报每隔一段时间Agent 调用一次服务健康检查工具发现状态码异常或者响应超时就自动生成一份故障描述并发给值班群。数据自动同步与清洗凌晨业务低峰期Agent 定时执行数据同步任务拉取上游数据文件调用清洗逻辑写入目标表再生成一份质量报告。个人助理类任务每天晚上提醒第二天的日程安排每周整理一次未读完的收藏文章这类任务更适合由轻量级 Agent 配合定时器来完成。这一类任务有一个共同特征单个任务执行时间可能不长但执行频率稳定、逻辑相对固定而且最好不需要人来盯着。这正是定时任务调度最擅长解决的问题。1.3 定时调度的几种实现路线对比在接 Agent-harness 之前我认真对比过三条技术路线方案适用场景优势短板操作系统 cron单机、简单定时脚本稳定、零依赖、通用无法感知 Python 环境异常重试和监控能力弱APSchedulerPython 项目内嵌的调度库进程内调度、支持多种触发器、可持久化任务单机单进程分布式场景需要额外设计Celery Beat分布式任务系统支持分布式执行、失败重试、任务队列需要额外部署消息队列对轻量场景偏重对于大多数 Agent-harness 项目来说我建议从 APScheduler 开始原因是它可以直接嵌在你的 Python 进程里不需要额外维护 Redis、RabbitMQ 这些基础设施也不需要去跟 crontab 的格式和环境变量纠缠。等以后任务量大了、需要多个 worker 协作时再考虑把执行层迁移到 Celery触发层保留 APScheduler 或者换成 Beat 也不迟。2. 动手前的准备工作安装 Agent-harness 与调度组件选型2.1 最小化运行环境搭建我这次的演示环境是 CentOS 7 的云服务器Python 3.10。无论你用的是虚拟环境还是 Docker建议先把 Python 版本锁定在 3.10 或以上因为后面操作时区、类型标注都会省心很多。mkdir agent-scheduler-demo cd agent-scheduler-demo python3 -m venv venv source venv/bin/activate pip install --upgrade pip pip install agent-harness安装完成后验证一下版本python -c import agent_harness; print(agent_harness.__version__)正常能打印出版本号说明框架核心装好了。如果你在安装过程中遇到依赖冲突大概率是 setuptools 或者 wheel 版本太低先升级这两个再装pip install --upgrade setuptools wheel随后安装定时调度需要的 APSchedulerpip install apscheduler这里多说一句APScheduler 3.x 和 4.x 的 API 有差异我教程里用的是 3.x 版本3.10.4如果你装到的是 4.x接口命名可能不一样建议直接固定版本安装避免文档和代码对不上。2.2 APScheduler 与 Agent-harness 的集成思路把调度器接进 Agent-harness最忌讳的是“调度器自己在跑Agent 自己在跑两边互不相通”。我在初学阶段犯过的错误就是——调度器到点执行了一个空函数函数内部再去手动调 Agent。虽然功能上也能实现但整个链路很别扭。更合理的思路是把 Agent 的启动入口设计成一个普通函数这个函数内部完成“创建 Agent - 编排任务 - 执行 - 回收结果”的完整流程调度器只负责按时间规则调用这个函数。这样一来调度器不关心 Agent 内部的实现细节Agent 也不依赖调度器的类型两边通过函数边界解耦后续要接 cron、要接 Celery 都只需要改一层包装。在代码组织上我建议至少拆出三个模块agent_runner.py负责把 Agent 执行流程封装成可调用函数。scheduler.py负责创建调度器、注册定时任务、启动事件循环。tasks/目录按业务域存放各种 Agent 任务的定义和编排逻辑。这样拆的好处是单个任务出问题不会影响其他任务调度逻辑和数据逻辑完全隔离。2.3 两个必须提前确认的关键前提第一个是进程必须常驻。APScheduler 是进程内调度它靠的是进程内部的事件循环不是操作系统级别的守护进程。如果你用nohup启动后终端一关进程就没了那调度肯定不生效。我在生产环境里用的是 systemd 服务来托管调度进程具体的 service 文件大概是这样[Unit] DescriptionAgent Scheduler Service Afternetwork.target [Service] Userwww-data WorkingDirectory/opt/agent-scheduler-demo ExecStart/opt/agent-scheduler-demo/venv/bin/python /opt/agent-scheduler-demo/scheduler.py Restartalways RestartSec5 [Install] WantedBymulti-user.target第二个是日志要单独落文件。调度任务很多时候是半夜执行的如果程序里的输出没有被正确记录第二天发现任务失败了你根本不知道凌晨三点发生了什么。我习惯在启动入口加一行 logging 配置import logging logging.basicConfig( filename/var/log/agent_scheduler.log, levellogging.INFO, format%(asctime)s - %(name)s - %(levelname)s - %(message)s )后面排查问题基本上全靠这份日志。3. 核心细节解析任务定义、注册与调度参数3.1 在 Agent-harness 中定义一个可被调度的 Agent 任务先不急着写调度得先把“Agent 任务本身”做成一个可以被反复调用的函数。拿我这次做的“每日行业动态简报”来说核心逻辑分三步资讯采集 Agent 从 RSS 源里抓取当天的新闻标题和链接。摘要 Agent 对每一条新闻生成一句话摘要。推送 Agent 把整理好的简报内容发送到钉钉群机器人。在task_briefing.py里我写了一个generate_daily_briefing(date_strNone)函数完整流程都封装在里面import asyncio from datetime import date from agent_harness import Agent, Workflow async def _run_briefing_pipeline(date_str: str) - dict: # 1. 创建采集 Agent collector Agent( namenews-collector, role资讯采集员, tools[fetch_rss, extract_content], backenddefault ) # 2. 创建摘要 Agent summarizer Agent( namenews-summarizer, role摘要提炼员, backenddefault ) # 3. 创建推送 Agent pusher Agent( namenotify-pusher, role消息推送员, tools[send_dingtalk_webhook], backenddefault ) # 编排工作流 workflow Workflow() workflow.add_stage(collect, collector, inputdate_str) workflow.add_stage(summarize, summarizer, depends_oncollect) workflow.add_stage(push, pusher, depends_onsummarize) result await workflow.run() return result def generate_daily_briefing(date_str: str None): 供调度器调用的任务入口。 APScheduler 无法直接调度 async 函数所以包一层同步函数。 if date_str is None: date_str date.today().isoformat() loop asyncio.new_event_loop() asyncio.set_event_loop(loop) try: result loop.run_until_complete(_run_briefing_pipeline(date_str)) logging.info(简报生成完成: %s, result) return result finally: loop.close()这里有一个非常关键的设计点APScheduler 的任务回调是同步函数没法直接await一个异步的 Agent 工作流。我刚开始写的时候直接把async函数注册进了调度器结果任务到点时只看到一句警告函数根本没有执行因为调度器把这个协程对象丢弃了。所以最终的调度入口必须是一个同步函数在函数内部再创建事件循环去跑 Agent 的协程。上面代码里asyncio.new_event_loop()那一段就是干这个事的。3.2 调度器配置触发器、任务存储与并发控制有了任务函数接下来就是往调度器里注册。APScheduler的BackgroundScheduler是最适合 Agent-harness 项目的调度器类型它不会阻塞主进程还会额外启动一个线程池来处理任务回调。from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger from apscheduler.jobstores.sqlite import SQLAlchemyJobStore from apscheduler.executors.pool import ThreadPoolExecutor jobstores { default: SQLAlchemyJobStore(urlsqlite:///jobs.db) } executors { default: ThreadPoolExecutor(max_workers10) } scheduler BackgroundScheduler(jobstoresjobstores, executorsexecutors) scheduler.add_job( generate_daily_briefing, triggerCronTrigger( hour8, minute0, day_of_weekmon-fri, timezoneAsia/Shanghai ), iddaily_briefing, replace_existingTrue, max_instances1, coalesceTrue, misfire_grace_time3600 ) scheduler.start()先解释几个新手容易忽视的参数。max_instances1表示同一个任务只允许一个实例在运行。这个参数格外重要因为有些 Agent 任务耗时比较长万一上一轮还没跑完下一轮触发条件又满足了调度器就会再拉一个新的实例形成并发。假如你用的是同一个 API Key 去调外部接口并发会导致限流假如任务里更新的是同一张数据库表就可能造成数据错乱。我经历过一次很尴尬的事故就是因为没设这个参数凌晨 2 点的清洗任务和凌晨 3 点的清洗任务重叠执行了同一条数据被两个实例同时更新最终结果完全乱了。coalesceTrue的作用是合并错过的执行。假设服务器因为维护关机了几个小时期间有多次调度任务本该运行但没运行开机后如果 coalesce 为 False调度器会把错过的所有任务一次性补执行如果 coalesce 为 True则只在恢复后执行最近一次。对大部分 Agent 任务来说补执行一堆“过期任务”没有意义我统一用的是 True。misfire_grace_time3600是用来兜底的。当任务因为主进程繁忙或者其他原因没能在预定时间启动只要错过的时长在宽限期内调度器仍然会补执行。3600 秒是 1 小时对于日报任务足够了。如果任务对时效性要求高比如每 5 分钟一次的状态检查这个值最好调小到 60120 秒。3.3 时区、任务 ID 与冲突控制三个必踩的坑第一坑时区不一致。我刚开始部署时调度进程跑在服务器上服务器时区是 UTC而我心里想的是北京时间 8 点结果任务每天下午 4 点才执行。排查了半天才发现是时区问题。通过 APScheduler 设置时区最稳妥的方式是在CronTrigger里显式声明timezone不要依赖操作系统时区。APScheduler 支持字符串形式的时区名CronTrigger( hour8, minute0, timezoneAsia/Shanghai )这个写法背后用的是标准时区库我建议你在设置完后先打印一下触发器对象确认它解析后的时区偏移量trigger CronTrigger(hour8, minute0, timezoneAsia/Shanghai) print(trigger) # 输出内容中应包含时区信息Asia/Shanghai 或 CST0800第二坑任务 ID 不稳定。add_job时必须显式指定稳定的id比如daily_briefing。如果你不指定APScheduler 每次都会生成一个随机 UUID。开发调试时还好但部署后你需要在持久化任务存储中管理这个任务比如临时暂停、修改触发规则如果 ID 是随机生成的就无从下手。而且当你需要更新某个任务的触发时间时务必设置replace_existingTrue否则新任务注册会失败直接在日志里报Job id (daily_briefing) conflicts with an existing job。第三坑执行锁缺失。max_instances1解决的是同一个调度器内部的并发问题。但如果你用了 systemd 托管又因为某些原因启动了多个调度进程每个进程都是独立的max_instances就不起作用了。这时候就需要引入跨进程锁。我不想引入 Redis 这类额外依赖就用最简单的方式在任务内部通过文件锁来保证同一时刻只有一个实例在跑。import fcntl import os LOCK_FILE /tmp/agent_daily_briefing.lock def generate_daily_briefing(date_str: str None): lock_fd open(LOCK_FILE, w) try: fcntl.flock(lock_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) except OSError: logging.warning(检测到任务已在其他进程运行本次跳过) lock_fd.close() return try: # ... 原有的 Agent 编排逻辑 pass finally: fcntl.flock(lock_fd, fcntl.LOCK_UN) lock_fd.close()LOCK_NB表示非阻塞模式如果拿不到锁说明另一个进程已经有实例在跑当前这次直接跳过。这个方案在单机多进程场景下足够可靠。4. 实操过程从零实现一个“每日简报 Agent”4.1 流程设计与子 Agent 拆分我选的案例是“每日行业动态简报 Agent”业务逻辑不复杂但能完整覆盖 Agent-harness 的编排调度场景。任务要求每天早上 8 点工作日自动抓取三个资讯源的最新内容筛选与“人工智能”相关的条目每条生成一句话摘要最后把 10 条以内的简报推送到钉钉工作群。我拆分成了三个子 Agent采集 Agent负责任务拆解调用 RSS 解析工具抓取指定资讯源的标题、链接、发布时间。摘要 Agent对采集结果逐条做摘要生成这一步本质上是调大语言模型接口但在 Agent-harness 里被封装成了 Agent 的一个动作。推送 Agent格式化输出通过钉钉 Webhook 机器人发送消息。这样的拆分不是为了炫技而是为了让每个 Agent 的职责单一、可单独测试。后续如果要把“摘要 Agent”换成不同的模型只需替换一个节点不影响采集和推送。4.2 调度入口与主程序实现还是回到代码。scheduler.py完整的内容如下import logging import asyncio from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger from apscheduler.jobstores.sqlite import SQLAlchemyJobStore from apscheduler.executors.pool import ThreadPoolExecutor from task_briefing import generate_daily_briefing logging.basicConfig( filename/var/log/agent_scheduler.log, levellogging.INFO, format%(asctime)s - %(name)s - %(levelname)s - %(message)s ) def main(): jobstores { default: SQLAlchemyJobStore(urlsqlite:///jobs.db) } executors { default: ThreadPoolExecutor(max_workers10) } scheduler BackgroundScheduler( jobstoresjobstores, executorsexecutors, timezoneAsia/Shanghai ) scheduler.add_job( generate_daily_briefing, triggerCronTrigger( hour8, minute0, day_of_weekmon-fri, timezoneAsia/Shanghai ), iddaily_briefing, replace_existingTrue, max_instances1, coalesceTrue, misfire_grace_time3600 ) scheduler.start() logging.info(调度器已启动等待任务触发...) try: asyncio.get_event_loop().run_forever() except (KeyboardInterrupt, SystemExit): scheduler.shutdown() if __name__ __main__: main()这里有个小细节scheduler.start()之后主程序不能直接退出否则调度进程就死了。我用的asyncio.get_event_loop().run_forever()让主线程一直挂起。如果你不想引入 asyncio用一个while True: time.sleep(60)的死循环也是可以的。总之进程必须常驻。启动时用nohup python scheduler.py /dev/null 21 启动后通过日志确认调度器注册成功tail -f /var/log/agent_scheduler.log正常应该看到类似“调度器已启动”的日志。为了进一步确认任务确实注册成功可以在代码里加一行打印scheduler.print_jobs()输出里会显示任务 ID、触发器表达式和下一次运行时间。看到“next run time”是北京时间 08:00:00就说明注册没问题。4.3 手动触发验证与真实运行结果等待到点再验证太慢了我习惯先手动触发一次确认 Agent 流程本身没有问题。在scheduler.py里临时加一行# 启动后立即手动触发一次验证整个链路 scheduler.add_job( generate_daily_briefing, triggerdate, run_datedatetime.now() timedelta(seconds5), idmanual_verify, replace_existingTrue )然后观察钉钉群有没有收到消息。这一步非常关键它把“调度触发”和“Agent 执行”分开验证了。如果手动触发能跑通就说明 Agent 编排没问题如果手动触发能跑通但定时触发不生效问题大概率出在调度配置上如果手动触发都跑不通那就要回去查 Agent 定义和工具调用。我的真实运行结果里采集 Agent 抓到了 23 条资讯摘要 Agent 逐条过滤后保留了 8 条跟目标主题强相关的推送 Agent 在 40 秒内完成了全部动作。整个过程没有人工介入。4.4 如何调整定时规则实用配置示例换一个需求你可能要改成“每天 23:30 运行一次”CronTrigger(hour23, minute30, timezoneAsia/Shanghai)改成“每周一、周三、周五早上 9 点 15 分”CronTrigger(day_of_weekmon,wed,fri, hour9, minute15, timezoneAsia/Shanghai)改成“每 30 分钟执行一次”from apscheduler.triggers.interval import IntervalTrigger IntervalTrigger(minutes30, timezoneAsia/Shanghai)改成“每月 1 号和 15 号早上 10 点”CronTrigger(day1,15, hour10, minute0, timezoneAsia/Shanghai)间隔触发和 cron 触发的区别在于IntervalTrigger是从任务启动时刻开始计时每 N 个单位时间执行一次CronTrigger则是严格按日历时间对齐。做“每小时整点执行”这类需求建议用CronTrigger(minute0)而不是IntervalTrigger(hours1)因为后者会随任务启动时间漂移。5. 常见问题与排查技巧实录5.1 任务到点没触发日志里什么都没留下这种问题是最让人头疼的因为“什么都没发生”本身就意味着信息缺失。我总结的排查顺序是先确认调度进程还在执行ps aux | grep scheduler.py看看进程是否存活。再确认下次运行时间在代码里加scheduler.print_jobs()检查任务的next run time是否在预期时间附近。如果显示的时区不是 0800说明时区配置有问题。查看系统日志如果进程突然被杀dmesg里可能有 OOM 记录或者 systemd 的Restart策略导致进程一直在重启。有一次我遇到的真实情况是服务部署到了新机器环境变量里没设置TZ程序逻辑里也没显式传时区APScheduler 就默认用了 UTC结果所有定时任务都比预期晚了 8 小时。从那以后我所有调度器的CronTrigger里都会强制写明timezoneAsia/Shanghai绝不把时区交给运行环境。5.2 任务重复执行同一条数据被处理两次这个前面提到过最常见的原因是max_instances没设成 1或者部署了多个调度进程。如果你用同一个 SQLite 的jobs.db文件却同时启动了多个scheduler.py进程两个进程会同时认为自己拥有这个任务到点后各跑一遍。文件锁可以兜底但更干净的方式是用 systemd 或者 supervisor 这类进程管理工具确保只存在一个调度进程。毕竟文件锁是一种补救手段单实例才是正解。5.3 重启后任务全部消失APScheduler 默认的内存任务存储只在进程运行期间有效进程一退出所有注册的任务就没了。解决方式是给调度器配置持久化 jobstore也就是我在前面代码里用的SQLAlchemyJobStore(urlsqlite:///jobs.db)。用了这个配置之后任务配置会被写入本地 SQLite 文件。scheduler 重启后会从 SQLite 里恢复任务但这里有个细节如果任务是通过scheduler.add_job()动态注册的恢复的任务依然在但如果你更新了任务代码老的任务记录可能还指向旧的函数名需要显式replace_existingTrue或者清空 jobstore 里的遗留任务否则会报TypeError: generate_daily_briefing() ...这类函数缺失错误。5.4 快速排查速查表现象可能原因解决方法任务到点无任何反应调度进程未运行 / 时区不对检查进程确认CronTrigger的 timezone 参数任务反复执行未设 max_instances 或存在多个调度进程max_instances1单实例部署重启后任务消失未配置持久化 jobstore使用 SQLite 的SQLAlchemyJobStore任务执行报“协程未运行”直接把 async 函数注册给了调度器包一层同步函数内部创建事件循环任务到点延迟很久才执行主进程阻塞 / 线程池耗尽调大ThreadPoolExecutor的 max_workers错过多次任务后一次性补跑coalesceFalse改成coalesceTrue5.5 日志排查的小技巧我最后再分享一个日志排查的经验。之前在/var/log/agent_scheduler.log里发现某次任务执行了 3 个小时还没结束但日志只在开头写了一句“任务开始”后面全无踪影。后来我把 Agent 内部的每一步编排都加了日志才定位到是外部资讯源接口卡住了HTTP 请求一直没设置超时。于是我在调用工具的代码里强制加了 15 秒超时问题立刻解决。给所有外部调用加超时这句话真的是老生常谈但它就是能避免大多数“任务卡死”类的疑难问题。Agent 任务往往要调多个外部服务任何一个服务不响应整个调度链都会卡住。建议每次在 Agent 的工具函数里都写清楚超时时间和重试策略好过上生产了再手忙脚乱地查日志。写在最后这次给 Agent-harness 接定时任务调度我个人的体会是调度的技术本身不难难的是把“定时”和“Agent 执行”这两件事完整接起来并考虑到异常、时区、并发、持久化这些边角问题。从一开始只会在 cron 里写一条命令到现在能清晰地把调度器、任务入口、Agent 编排拆成不同的层次过程中的收获不仅仅是多会了一个工具更多的是对整个任务生命周期有了把控感。如果你接下来也要在 Agent-harness 里接定时任务我建议先从一个低频、非核心的任务开始跑一周把日志、监控、异常处理都摸通以后再让调度器去承载那些真正重要的业务任务。定时调度就像一个靠谱的同事你交代清楚规则它会一直替你盯着时间但前提是你得先把“后路”都铺好它才不会在关键时刻掉链子。
返回列表