跳转至

Scheduler API

Scheduler(config=None, backend=None, logger=None, enable_metrics=False, use_global_scheduler=True)

高性能异步任务调度器。

特性: - 支持秒级精度的定时任务 - 可插拔的后端存储(内存/Redis/RabbitMQ/SQLite) - 可插拔的日志系统(structlog/loguru/stdlib) - 优雅关闭和信号处理 - 自动重试机制 - 任务优先级和并发控制

初始化调度器。

参数

config: 调度器配置,默认使用 SchedulerConfig() backend: 队列后端,默认依据配置自动实例化 logger: 日志适配器,默认使用 structlog enable_metrics: 是否启用性能指标收集 use_global_scheduler: 是否将当前实例注册为全局调度器,以支持装饰器自动注册

set_logger(logger)

设置日志适配器。

参数

logger: 日志适配器实例

示例

from loguru import logger from chronflow.logging import LoguruAdapter

scheduler.set_logger(LoguruAdapter(logger))

register_task(task)

注册任务到调度器。

参数

task: 要注册的任务实例

抛出

ValueError: 任务名称已存在

unregister_task(task_name)

注销任务。

参数

task_name: 任务名称

get_task(task_name)

获取任务实例。

参数

task_name: 任务名称

返回值

任务实例,如果不存在返回 None

start(daemon=False) async

启动调度器。

参数

daemon: 是否以守护进程模式运行

返回值

守护进程模式下返回子进程 PID,否则返回 None

stop(daemon=False, *, pid=None, name=None, timeout=None) async

停止调度器或守护进程。

参数

daemon: 是否操作守护进程 pid: 指定守护进程 PID name: 指定守护进程名称 timeout: 等待终止的超时时间

restart(daemon=False, *, pid=None, name=None, timeout=None) async

重启调度器或守护进程。

cleanup(daemon=False, *, pid=None, name=None) async

清理守护进程僵尸状态。

stop_daemon(*, pid=None, name=None, timeout=None) async

保持兼容的守护进程停止接口。

restart_daemon(*, pid=None, name=None, timeout=None) async

保持兼容的守护进程重启接口。

cleanup_daemon(*, pid=None, name=None) async

保持兼容的守护进程清理接口。

run_context() async

使用上下文管理器运行调度器。

示例

async with scheduler.run_context(): # 调度器在这里运行 await asyncio.sleep(60)

自动停止调度器

get_stats() async

获取调度器统计信息。

返回值

包含统计信息的字典

list_tasks()

获取所有任务列表。

返回值

任务信息列表

示例

tasks = scheduler.list_tasks() for task_info in tasks: print(f"{task_info['name']}: {task_info['status']}")

get_task_count()

获取各状态任务数量统计。

返回值

任务数量统计字典

示例

counts = scheduler.get_task_count() print(f"运行中: {counts['running']}") print(f"失败: {counts['failed']}")

get_task_by_status(status)

根据状态获取任务列表。

参数

status: 任务状态

返回值

符合状态的任务列表

示例

failed_tasks = scheduler.get_task_by_status(TaskStatus.FAILED) for task in failed_tasks: print(f"失败任务: {task.config.name}")

get_task_by_tag(tag)

根据标签获取任务列表。

参数

tag: 标签名称

返回值

包含该标签的任务列表

示例

critical_tasks = scheduler.get_task_by_tag("critical")

pause_task(task_name) async

暂停任务(禁用)。

参数

task_name: 任务名称

返回值

是否成功暂停

示例

await scheduler.pause_task("my_task")

resume_task(task_name) async

恢复任务(启用)。

参数

task_name: 任务名称

返回值

是否成功恢复

示例

await scheduler.resume_task("my_task")

get_metrics()

获取性能指标。

返回值

性能指标字典,如果未启用指标收集则返回 None

示例

if scheduler.metrics_collector: metrics = scheduler.get_metrics() print(f"总执行次数: {metrics['total_executions']}") print(f"成功率: {metrics['success_rate']:.2%}")

export_prometheus_metrics()

导出 Prometheus 格式的指标。

返回值

Prometheus 格式的指标文本,如果未启用指标收集则返回 None

示例

metrics_text = scheduler.export_prometheus_metrics() if metrics_text: # 可以通过 HTTP 端点暴露给 Prometheus print(metrics_text)

reset_metrics()

重置性能指标。

示例

scheduler.reset_metrics()

discover_tasks_from_directory(directory, *, pattern='task.py', recursive=True, exclude_patterns=None)

从目录自动发现并注册任务。

参数

directory: 要扫描的目录路径 pattern: 文件名匹配模式,支持通配符 (默认: "task.py") recursive: 是否递归扫描子目录 (默认: True) exclude_patterns: 排除的文件名模式列表

返回值

发现并注册的任务列表

示例

扫描所有 task.py 文件

scheduler.discover_tasks_from_directory("my_app/modules")

扫描所有 *_tasks.py 文件

scheduler.discover_tasks_from_directory( "my_app", pattern="tasks.py", exclude_patterns=["test.py"] )

discover_tasks_from_package(package_name, *, pattern='task.py', exclude_patterns=None)

从包自动发现并注册任务。

参数

package_name: 包名 (例如: "my_app.tasks") pattern: 文件名匹配模式 exclude_patterns: 排除的文件名模式列表

返回值

发现并注册的任务列表

示例

扫描包及其子包

scheduler.discover_tasks_from_package("my_app.tasks")

扫描特定模式的文件

scheduler.discover_tasks_from_package( "my_app", pattern="*_tasks.py" )

discover_tasks_from_modules(module_names)

从指定的模块列表中发现并注册任务。

参数

module_names: 模块名列表 (例如: ["app.tasks.user", "app.tasks.email"])

返回值

发现并注册的任务列表

示例

scheduler.discover_tasks_from_modules([ "my_app.tasks.user_tasks", "my_app.tasks.email_tasks", ])

__repr__()

字符串表示。