Skip to main content

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.filesfiles 是 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 必须导出名为 tasksTaskRegistry。包内可以继续使用 相对导入拆分多个任务模块。管理员通过受信网络发布 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

scheduler_sdk-0.2.0.tar.gz (53.1 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

scheduler_sdk-0.2.0-py3-none-any.whl (23.3 kB view details)

Uploaded Python 3

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

Hashes for scheduler_sdk-0.2.0.tar.gz
Algorithm Hash digest
SHA256 eda79d5abe44d5b10bd3ac315bdb1d4938c1b3ff21e04f7028fd4d47e1f728e3
MD5 4e8fa3e047da59e945cf11b8b0eaef40
BLAKE2b-256 f352b7a715b84a5f10005e744a26c839b5aa5bbdc14fa2080a0f138bf40ae910

See more details on using hashes here.

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

Hashes for scheduler_sdk-0.2.0-py3-none-any.whl
Algorithm Hash digest
SHA256 3cfa49338f291a2391800adf84fdb82d87006ce3db96c31950caf2c734701642
MD5 84c499c8071edbe2ec764df11c7df434
BLAKE2b-256 8a077c7f94c833ef38af27c239957d880373e5ca846b7357729df0728ae79fef

See more details on using hashes here.

Release history Release notifications | RSS feed

0.2.2

2 files

0.2.1

2 files

This release

0.2.0 This release

2 files

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page