
Swarms CronJob 定时任务编排实战从单 Agent 周期调度到多 Agent 混合节奏运维【免费下载链接】swarmsThe Enterprise-Grade Multi-Agent Orchestration Framework. Website: https://swarms.ai项目地址: https://gitcode.com/GitHub_Trending/swar/swarms导读本文将围绕 Swarms 框架内置的CronJob调度器系统讲解如何把任意Agent甚至整个ConcurrentWorkflow绑定到指定时间间隔上自动周期执行。从一个 Agent 一个间隔的最小可用任务到run_many管理多 Agent 各自独立节奏的任务舰队再到非阻塞控制、回调加工输出、错误预算与实时监控你将掌握一套可直接用于生产环境的定时 Agent 任务编排方案。文中所有结论均以 examples/guides/deployment/cron_job_examples 目录下的 13 个示例、核心实现 cron_job.py 及测试 test_cron_job.py 为依据。CronJob 核心模型一个 Agent 绑定一个间隔CronJob的核心设计约束非常明确一个CronJob实例只把一个 Agent或任意可调用对象绑定到一个时间间隔上并持续运行直到你主动停止它。这一点在源码类注释中被反复强调cron_job.pyOneCronJobbindsoneagent tooneinterval. It schedules the task, runs it on its own background thread, and keeps running it until you stop it.它的工作方式是run(task)把任务注册进调度器随后阻塞调用线程Agent 在后台线程里按间隔被反复执行通过Ctrl-CKeyboardInterrupt或从其他线程调用stop()结束运行。间隔格式数字单位间隔采用numberunit字符串格式单位支持秒、分、时三种单复数均可。合法示例30seconds—— 每 30 秒10minutes—— 每 10 分钟1hour—— 每小时源码中的_parse_intervalcron_job.py用正则(\d)(\w)解析数字与单位并映射到schedule库的every(x).seconds / .minutes / .hours方法上。以下约束在初始化时reliability_check即被强校验非法配置会抛出CronJobConfigError数字必须大于 00second会调度出一个永不触发的任务源码中专门注释说明了这一点单位仅限second(s)、minute(s)、hour(s)1day、1 second、second等均会被拒绝间隔不能为空字符串或纯空白。测试 test_cron_job.py 分别对合法间隔1second、5seconds、1minute、10minutes、1hour、2hours和非法间隔0second、0minutes、空串、-1second、1day、second、1 second做了参数化验证。五个入门示例从最小任务到生产级韧性1. 最小可用任务single_agent_cron.pysingle_agent_cron.py 是最小的有用 CronJob——一个 Agent、一个间隔运行直到你停止它from swarms import Agent, CronJob agent Agent( agent_nameMarket-Watcher, agent_descriptionWatches the market and reports anything notable, system_prompt( You are a market analyst. Report only what is notable since the last check. Three bullets maximum. If nothing is notable, say so in one line rather than padding. ), model_namegpt-5.4, max_loops1, print_onTrue, ) job CronJob(agentagent, interval30seconds) if __name__ __main__: print(Running every 30 seconds. Ctrl-C to stop.\n) job.run( Summarise anything notable in the AI chip market right now. )要点run()阻塞调用线程Agent 在后台线程按间隔执行见start()创建的cronjob_{job_id}线程cron_job.py每个 tick 传入的同一个任务字符串会原样交给 Agent建议在__main__中运行方便用 Ctrl-C 优雅退出。2. 一个 Agent 多个任务共享节奏one_agent_many_tasks_cron.py当多个检查项共享同一节奏、且由同一个 Agent 完成时使用batched_run。它把任务列表里的每一项都在该作业的间隔上逐一注册因此每个 tick 都会执行全部任务one_agent_many_tasks_cron.pyCHECKS [ Check whether inventory levels are below reorder thresholds., Check whether refund volume is above its weekly average., Check whether any support queue has waited longer than an hour., ] if __name__ __main__: print(f{len(CHECKS)} checks every 15 minutes. Ctrl-C to stop.\n) CronJob(agentagent, interval15minutes).batched_run(CHECKS)从实现看cron_job.pybatched_run会先循环调用_run(task)把每个任务都注册进调度器然后才_block_forever()阻塞等待。测试 test_cron_job.py 专门验证了两个行为batched_run会把所有任务都调度上曾经存在阻塞在第一个任务、其余任务未调度的 bugtest_batched_run_schedules_every_task即回归测试它按输入顺序返回每个任务对应的一个schedule.Job数量与任务数一致。3. 多 Agent 各自节奏multi_agent_schedules_cron.py一个 CronJob 绑定一个 Agent 到一个间隔意味着混合节奏的任务舰队需要每个 Agent 一个 Job。CronJob.run_many批量构建、同时启动、并一次性阻塞multi_agent_schedules_cron.pyCronJob.run_many( [ { agent: price_agent, interval: 30seconds, task: Check the BTC price., }, { agent: anomaly_agent, interval: 10minutes, task: Scan recent market data for anomalies., # This one talks to a flaky data source, so give it a budget: # five failures in a row and it stops rather than retrying # forever. The other two are unaffected either way. max_consecutive_errors: 5, }, { agent: digest_agent, interval: 1hour, task: Summarise the last hour of market activity., }, ] )每份 schedule 字典中agent、interval、task为必填键job_id、callback、max_consecutive_errors、kwargs会透传给 Agent 的run为可选键。run_many对空列表、缺失必填键都会抛出CronJobConfigError并指明具体缺失项cron_job.py测试见 test_cron_job.py。隔离性是关键特性每个 Job 拥有独立的调度线程慢的或失败的 Agent不会拖延或停止其他 Agent且各自携带独立的错误计数和错误预算。阻塞模式下run_many会等待所有 Job 结束期间 Ctrl-C 会触发stop_many若存在耗尽错误预算而停止的 Job最终会抛出CronJobExecutionError并汇总每个失败 Job 的连续失败次数与最后错误cron_job.py。4. 非阻塞启动与运行中巡检non_blocking_cron.py当调度不是进程的主业比如 Web 服务器、Bot、Notebook 场景时用run_many(..., blockFalse)让作业启动后立即返回把 Job 对象交还给你non_blocking_cron.pyjobs CronJob.run_many( [ {agent: agents[fast], interval: 2seconds, task: poll, job_id: fast-poller}, {agent: agents[medium], interval: 5seconds, task: check, job_id: medium-checker}, {agent: agents[slow], interval: 10seconds, task: summarise, job_id: slow-summariser}, ], blockFalse, # start them, hand control back ) print(Three jobs started. Doing other work for 20 seconds...\n) time.sleep(20) print(\nMid-flight status:) for job in jobs: stats job.get_execution_stats() print( f {stats[job_id]:18} every {stats[interval]:10} fok{stats[execution_count]:3} ffailed{stats[error_count]:3} fup{stats[uptime]:.0f}s ) CronJob.stop_many(jobs)生命周期由你接管用get_execution_stats()巡检用CronJob.stop_many(jobs)统一关闭。这个示例使用了一个极简的EchoAgent仅暴露run(task, **kwargs)替代真实 Agent无需任何 API Key 即可运行是验证 CronJob 机制本身的最佳起点。5. 韧性设计resilient_cron.py长时间运行的作业必然遇到限流、断连、模型服务抖动。CronJob 的失败处理哲学与系统 cron 一致任务抛错会被记录并在下一个 tick 重试而不是让整个调度崩溃。两个核心控制点resilient_cron.pymax_consecutive_errors设置连续失败多少次后停止作业而不是无限重试。不设置默认None则永不放弃。当预算耗尽时作业停止且run()会抛出CronJobExecutionError从而保证一个已死的调度绝不会被误认为健康get_execution_stats()在作业运行期间从其他线程轮询查看成功数、失败数与最后一次错误。job CronJob( agentFlakyAgent(), # 约一半概率抛 ConnectionError interval2seconds, max_consecutive_errors10, # 连续 10 次失败说明真的出问题了 ) threading.Thread(targetmonitor, args(job,), daemonTrue).start() threading.Timer(60, job.stop).start() # 一分钟后自动停止示例可自行退出 try: job.run(Fetch the latest reading.) except CronJobExecutionError as e: print(f\nJob gave up: {e}) # 只有预算耗尽才会走到这里 else: stats job.get_execution_stats() # 正常 stop() 则正常返回 print(f\nStopped cleanly after {stats[execution_count]} successful run(s) and {stats[error_count]} failure(s).)示例中的monitor线程在启动前会先自旋等待job.is_running变为True因为监控线程先于job.run()启动此时is_running还是False直接循环会立即退出这是一个值得复用的观察技巧。从源码层面看调度循环_run_schedulecron_job.py在每次schedule.run_pending()成功后把consecutive_errors归零捕获到异常时累加error_count与consecutive_errors并记录last_error只有当连续失败数达到max_consecutive_errors时才置_stopped_due_to_error True并退出循环。run()在阻塞结束后调用_raise_if_stopped_by_errors()cron_job.py区分干净停止与因错误放弃两种状态。测试对这三条语义做了逐条验证test_cron_job.pytest_a_failed_execution_does_not_kill_the_schedule一次瞬时失败后调度仍存活、计数器复位test_consecutive_errors_stop_the_job_when_a_budget_is_set设置预算后连续失败会抛CronJobExecutionError并置位_stopped_due_to_errortest_without_a_budget_the_job_retries_indefinitely不设预算则持续重试、调度线程存活test_stats_expose_the_failure_pictureget_execution_stats()能如实暴露error_count、consecutive_errors、last_error。回调Callback运行期输出加工CronJob构造函数和run_many的 schedule 字典都支持callback参数。回调签名固定为callback(output, task, metadata) - Any其中outputAgent 的原始输出task本次执行的任务字符串metadata包含job_id、timestamp、execution_count、task、kwargs、start_time、is_running的字典。回调的返回值会成为该次 tick 的最终产出源码见 cron_job.py。若回调自身抛异常会被记录日志并退回使用原始输出——回调故障绝不影响作业继续运行。set_callback(callback)还支持在作业创建后动态更新回调cron_job.py。基础用法simple_callback_example.pysimple_callback_example.py 展示了最朴素的模式——给每次输出附加执行元数据def simple_callback(output, task, metadata): return { agent_output: output, execution_number: metadata[execution_count], timestamp: datetime.fromtimestamp(metadata[timestamp]).isoformat(), task: task, job_id: metadata[job_id], } cron_job CronJob( agentagent, interval10seconds, job_idsimple-callback-example, callbacksimple_callback, )四种进阶模式callback_cron_example.pycallback_cron_example.py 系统演示了四类回调应用输出结构化转换transform_output_callback把 Agent 输出包装为包含original_output、transformed_at、execution_number、task_executed、job_status、uptime_seconds的结构化字典过滤与增强filter_and_enhance_callback按关键词important、key、significant、trend判定输出优先级并附加priority、analysis_type等字段实时监控类MonitoringCallback以可调用类实现维护最近 100 条输出历史累计success_count/error_count计算成功率、单次执行耗时并提供get_summary()汇总API 集成api_webhook_callback把输出打包成api_payload含job_id、execution_id、timestamp后外发示例中以日志代替真实 HTTP 请求。该示例还为四个回调各创建了一个 Job间隔分别 15s/20s/25s/30s在独立线程中并发启动运行 2 分钟后打印监控汇总与各 Job 的get_execution_stats()Ctrl-C 时统一stop()并输出最终监控总结。领域实战行情与股票的定时分析README 中的Callbacks and domain examples组集中展示了把 CronJob 应用于真实金融数据场景的三种完整套路。股票监控cron_job_example.py 与 figma_stock_example.pycron_job_example.py 定义了get_figma_stock_data(stock)通过 Yahoo Financeyfinance抓取一只股票的全维度数据组装为 JSON 字符串涵盖company_info公司名、代码、行业、官网、业务描述current_market_data现价、昨收、开盘、日内高低、成交量、市值、涨跌额与涨跌幅financial_metricsPE、前瞻 PE、市净率、市销率、企业价值、Beta、股息率、派息率trading_statistics50 日均线、200 日均线、52 周高低、流通股本等recent_performance近 30 天起止价、总收益率、最高最低价、平均成交量real_time_datafast_info的实时价、量、买卖盘口。文件底部以注释形式保留了完整的量化交易 Agent CronJob用法把get_figma_stock_data(FIG)注入系统提示词构造带dynamic_temperature_enabledTrue、output_typestr-all-except-first等参数的 Agent再以CronJob(agentagent, interval10seconds).run(task...)驱动。figma_stock_example.py 则演示消费该数据函数打印完整 JSON、提取现价/市值/PE 三个关键指标并根据当日涨跌输出带 emoji 的简报。swarms_tools 集成cron_job_figma_stock_swarms_tools_example.pycron_job_figma_stock_swarms_tools_example.py 展示另一条数据获取路径——直接使用swarms_tools包内置的yahoo_finance_api([FIG])并给出三种方法的对比直接调用swarms_tools.yahoo_finance_api拿原始数据把yahoo_finance_api作为tools[yahoo_finance_api]注入 Agent配合 FINANCIAL_AGENT_SYS_PROMPT 金融系统提示词让 Agent 自主取数并分析与自定义get_figma_stock_data_simple()的输出做比对。这说明 CronJob 的 Agent 既可以是自带工具取数的智能体也可以是先取数再注入提示词的瘦 Agent两种形态都可被定时驱动。加密货币并发分析simple_concurrent_crypto_cron.py 与 crypto_concurrent_cron_example.py两个示例展示一个高价值模式CronJob 调度的不一定是单个 Agent而可以是整个ConcurrentWorkflow。调度器按agent.run/agent(task)的鸭子类型分发cron_job.py优先使用run()方法否则把 Agent 当普通可调用对象直接调用因此ConcurrentWorkflow天然满足该接口。simple_concurrent_crypto_cron.py 是最简版本四步走为 BTC / ETH / SOL 各建一个只分析自己负责币种的专家 Agentprint_onFalse便于并发输出组装ConcurrentWorkflow(name..., agentsagents, show_dashboardTrue)用CronJob(agentworkflow, interval60seconds)包装把 Workflow 当作 Agent通过 CoinGecko API 取实时价格数据拼进任务提示词cron_job.run(tasktask)每分钟并发分析一次。crypto_concurrent_cron_example.py 是完整版六个专家 AgentBTC、ETH、SOL、ADA、BNB、XRP每个都带coin_gecko_coin_api工具并配深度专业提示词通过ConcurrentWorkflow并行执行CronJob每 5 秒驱动一轮机构级加密货币分析架构注释CronJob - ConcurrentWorkflow - [Bitcoin Agent, Ethereum Agent, ...] - Parallel Analysis。Solana 价格追踪器solana_price_tracker.pysolana_price_tracker.py 是一个端到端完整范例get_solana_price()请求 CoinGecko/simple/price接口带 10 秒超时取 SOL 的 USD 价格、市值、24h 成交量与涨跌幅并做格式化analyze_solana_data(data)基于 24h 涨跌幅生成带情感倾向的分析文本涨超 5% 为强看多、跌超 5% 为强看空并结合成交量/市值比判断市场活跃度把get_solana_price()的实时数据注入system_promptCronJob(agentagent, interval30seconds).run(taskAnalyze the current Solana (SOL) price data comprehensively...)定时产出完整市场报告。这是数据获取函数 提示词注入 CronJob 定时驱动三件套的完整落地可作为任何周期性数据监控任务的模板。异常体系与统计字段速查CronJob定义了三级异常cron_job.py异常触发场景CronJobError所有 CronJob 异常基类CronJobConfigErrorAgent 缺失、interval 为空/为零/格式非法、run_many传入空列表或缺必填键CronJobExecutionError调度失败、任务执行失败或作业因耗尽max_consecutive_errors而放弃消息会指明失败次数与最后错误get_execution_stats()返回的字段cron_job.py字段含义job_id作业唯一标识未指定时自动生成is_running调度线程是否存活execution_count成功执行次数error_count累计失败次数consecutive_errors自上次成功以来的连续失败数last_error最近一次失败异常字符串或Nonestopped_due_to_error仅在耗尽错误预算而放弃时为Truestart_time/uptime启动时间戳与已运行秒数interval间隔字符串从示例到生产工程建议综合 README 的 Overview 与示例实现可提炼如下实践准则按节奏切分 Job共享节奏的多任务用batched_run不同节奏的多个 Agent 用run_many避免为每个任务单独建线程善用隔离性run_many的每 Job 独立线程与独立错误预算意味着一个 Agent 的限流或崩溃不会拖垮整个舰队可放心混排不同可靠性的任务永远给不可靠数据源设错误预算如multi_agent_schedules_cron.py中给对接不稳定数据源的 Job 设置max_consecutive_errors配合CronJobExecutionError让死调度显式暴露把监控做进回调MonitoringCallback这类可调用类能在作业存活期内累积成功率、执行耗时与历史输出是轻量级可观测性方案阻塞与非阻塞按场景选CLI 工具用阻塞run()/run_many()加 Ctrl-C服务型进程用blockFalseget_execution_stats()巡检 stop_many()接管生命周期用无 Key 假 Agent 先验证机制non_blocking_cron.py与测试中的MockAgent表明只要对象暴露run()或可调用即可被调度调试调度逻辑不必消耗真实模型调用。若需要将这类定时任务进一步接入 Web 服务或外部编排环境可参考 deployment 目录下 FastAPI 与 cron 集成相关的更多部署文档。【免费下载链接】swarmsThe Enterprise-Grade Multi-Agent Orchestration Framework. Website: https://swarms.ai项目地址: https://gitcode.com/GitHub_Trending/swar/swarms创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考