Scheduler Worker SDK
当前实现包含任务声明与 Worker Runtime。一个脚本可以组合并运行多个任务函数,
基础设施配置全部来自当前目录的 .env。
安装
发布包是纯 Python 通用 wheel,同时支持 Windows 和 Linux:
python -m pip install scheduler-sdk==0.2.0
需要 Python 3.10 或更高版本。
import logging
from scheduler_sdk import TaskExecution, TaskRegistry, Worker
logging.basicConfig(level=logging.INFO)
tasks = TaskRegistry("assets")
@tasks.task("report.import")
async def import_report(
task_context: str,
execution: TaskExecution,
overwrite: bool = False,
) -> dict:
source = execution.files[0].path
return {"source": str(source), "overwrite": overwrite}
Worker(tasks).run()
也可以使用 tasks.register(name, function) 接入已有函数,使用
tasks.include(other_registry) 显式组合多个业务模块。固定任务字段按同名参数注入,
kv_arg 按 Python 关键字参数规则展开;SDK 不执行隐式类型转换。
运行配置
在 Worker 运行目录创建 .env;源码包内提供 .env.example 模板,支持以下参数:
NATS_URL:NATS 内网地址,必填。NATS_TOKEN:NATS 令牌,可留空。WORKER_CLIENT_ID:Worker 唯一标识,必填。WORKER_CHANNEL:Worker 通道,默认default。WORKER_CONCURRENCY:进程内全局并发数,默认4。WORKER_HEARTBEAT_INTERVAL:心跳间隔秒数,默认10。SCHEDULER_HTTP_URL:任务包服务完整 HTTP 根地址;不配置则关闭自更新。WORKER_FILE_HTTP_TIMEOUT:任务文件下载超时秒数,默认30。WORKER_FILE_MAX_BYTES:单个任务文件下载上限,默认100MB。WORKER_HEALTH_FILE:Daemon 读取的健康文件路径,默认使用系统临时目录。WORKER_HEALTH_INTERVAL:健康文件刷新间隔秒数,默认5。WORKER_BUNDLE_ROOT:本地任务包状态目录,默认.scheduler-worker。WORKER_BUNDLE_UPDATE_INTERVAL:运行中检查间隔秒数,默认60。WORKER_BUNDLE_HTTP_TIMEOUT:HTTP 超时秒数,默认30。WORKER_BUNDLE_MAX_BYTES:任务包 ZIP 最大字节数,默认20MB。
进程启动后,任务注册表会在连接 NATS 前冻结。同步任务在线程中执行,异步任务直接
运行在事件循环中,两者共享同一个全局并发限制。SDK 不设置任务执行超时,也不读取
kv_arg.timeout;同步函数真正结束前始终占用并发槽。
任务文件
服务端通过 multipart 接收文件后,会在任务消息中注入文件元数据。配置
SCHEDULER_HTTP_URL 后,SDK 在调用任务函数前下载并核对文件大小,将本地临时路径
放入 TaskExecution.files。files 是 SDK 保留输入,不会作为普通 kv_arg 参数展开。
from scheduler_sdk import TaskExecution
@tasks.task("report.import")
def import_report(execution: TaskExecution) -> dict:
task_file = execution.files[0]
return {
"name": task_file.name,
"content": task_file.path.read_text(encoding="utf-8"),
}
本地文件仅在任务函数执行期间有效,函数返回或抛出异常后 SDK 会清理任务临时目录。
缺少 HTTP 地址、下载失败、超过大小上限或下载大小与元数据不一致时,SDK 返回
task_file_* 结构化错误,并且不会调用任务函数。
错误返回
任务函数可以抛出 TaskError 返回可判断的业务错误:
from scheduler_sdk import TaskError
raise TaskError(
"报表格式不支持",
code="report_invalid",
details={"line": 3},
)
TaskError 会返回 code/message/details。其他异常会返回异常类型与消息,完整 traceback
只写入 Worker 日志,不通过 NATS 返回。
任务包
任务包使用固定结构,不携带或安装 Python 依赖:
tasks.zip
├── bundle.json
└── task_bundle/
├── __init__.py
└── ...
bundle.json 只要求非空版本号:
{"version": "2026.07.27.1"}
task_bundle/__init__.py 必须导出名为 tasks 的 TaskRegistry。包内可以继续使用
相对导入拆分多个任务模块。管理员通过受信网络发布 ZIP 并同时更新通道目标:
curl -F "bundle=@tasks.zip" \
http://192.0.2.10:17002/example/scheduler/admin/task-bundles/stable
Worker 启动前会下载目标包、核对 SHA-256、安全解压并验证 Python 入口,然后原子更新 激活指针。运行中发现新版本时,后台只完成下载与结构校验,不执行新任务函数代码;随后 停止接收新任务、等待在途任务结束并重新执行当前脚本。当前版本无法加载时自动回滚上一 版本,本地只保留当前与上一版本。任务包状态与错误原因会通过 Heartbeat v2 上报。
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file scheduler_sdk-0.2.0.tar.gz.
File metadata
- Download URL: scheduler_sdk-0.2.0.tar.gz
- Upload date:
- Size: 53.1 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via:
uv/0.11.26 {"installer":{"name":"uv","version":"0.11.26","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"22.04","id":"jammy","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
eda79d5abe44d5b10bd3ac315bdb1d4938c1b3ff21e04f7028fd4d47e1f728e3
|
|
| MD5 |
4e8fa3e047da59e945cf11b8b0eaef40
|
|
| BLAKE2b-256 |
f352b7a715b84a5f10005e744a26c839b5aa5bbdc14fa2080a0f138bf40ae910
|
File details
Details for the file scheduler_sdk-0.2.0-py3-none-any.whl.
File metadata
- Download URL: scheduler_sdk-0.2.0-py3-none-any.whl
- Upload date:
- Size: 23.3 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via:
uv/0.11.26 {"installer":{"name":"uv","version":"0.11.26","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"22.04","id":"jammy","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
3cfa49338f291a2391800adf84fdb82d87006ce3db96c31950caf2c734701642
|
|
| MD5 |
84c499c8071edbe2ec764df11c7df434
|
|
| BLAKE2b-256 |
8a077c7f94c833ef38af27c239957d880373e5ca846b7357729df0728ae79fef
|