统一的计算节点框架 — 调度器(Scheduler)与执行器(Executor)于一体,支持水平扩展、高可用的任务分发与执行
Project description
lzm-workers — 统一的计算节点框架
编码:utf8 | 作者:Lzm | 日期:2026-07-16
调度与执行彻底分离 — lzm-workers 一个包,两种模式:
- 调度器(Scheduler):任务管理、分发、定时、重试、统计
- 执行器(Executor):任务消费、插件/脚本执行、结果回传
支持水平扩展、高可用部署,通过 lzm-edsm 事件驱动通信,实现零轮询、高吞吐、低延迟的任务分发。
安装
# 核心依赖
pip install lzm-workers
# 带 MySQL 支持(调度器持久化)
pip install lzm-workers[mysql]
# 带 EDSM 传输
pip install lzm-workers[edsm-redis] # Redis 传输
pip install lzm-workers[edsm-rabbitmq] # RabbitMQ 传输
pip install lzm-workers[edsm-all] # 全部传输
# 全量
pip install lzm-workers[all]
快速开始
方式一:CLI 启动
# 启动执行器守护进程
lzm-workers executor --config config.yaml
# 启动调度器
lzm-workers scheduler --config config.yaml
方式二:嵌入代码
调度器模式 — 可嵌入任意 Python 服务:
import asyncio
from lzm.workers import Scheduler
from lzm.workers.common.config import Settings
async def main():
scheduler = Scheduler(Settings(mode="scheduler"))
await scheduler.start()
# 创建一个插件任务
record = await scheduler.create_task(
task_type="plugin",
task_name="heading_strict",
payload={"text": "..."},
priority=5,
timeout=300,
)
print(f"任务已创建: {record.task_id}")
# 保持运行
await asyncio.Event().wait()
asyncio.run(main())
执行器模式 — 独立守护进程:
import asyncio
from lzm.workers import ExecutorRunner
from lzm.workers.common.config import Settings
async def main():
runner = ExecutorRunner(Settings(mode="executor"))
await runner.run_forever()
asyncio.run(main())
核心架构
┌── lzm-workers 调度器模式 (可嵌入任意服务) ──────────────────────┐
│ TaskManager → Dispatcher → Reaper → EDSM Event → Redis │
└──────────────────────────────┬───────────────────────────────────┘
│
EDSM Event (COMPETING)
Redis 心跳 (SETEX)
│
┌── lzm-workers 执行器模式 (独立守护进程) ────────────────────────┐
│ EDSM Consumer → TaskRunner(租约) → ExecutorPool → 结果回写 │
│ HeartbeatReporter (每 3s) → HealthServer (/health, /metrics) │
└──────────────────────────────────────────────────────────────────┘
3 道健壮防线
| 防线 | 机制 | 响应时间 |
|---|---|---|
| 租约 | Redis SET NX EX task:lease:{id} |
任务超时 + 30s 后自动释放 |
| 心跳 + Reaper | 每 3s 刷新 TTL=10s 心跳,Reaper 每 30s 巡检 | ~30s 回收失联执行器 |
| 调度器无状态 | 状态存 MySQL + Redis,重启自动恢复 | 瞬时恢复 |
性能特性
- 零轮询:纯事件驱动,无任何轮询操作
- 批量回写:结果攒 50 条或每 3 秒批量 flush
- Redis O(1):心跳、租约、发现均为 Redis O(1) 操作
- 异步并发:asyncio.Semaphore 控制并发,非阻塞调度
使用场景
| 场景 | 推荐模式 | 说明 |
|---|---|---|
| Web 服务需要异步执行任务 | 调度器嵌入 Web 服务 | 请求进来 → 创建任务 → 立即返回 → 异步执行 |
| 独立计算节点集群 | 执行器守护进程 | 多台机器各自运行 Executor,水平扩展 |
| AI 流水线处理 | 调度器 + 执行器 | 调度器编排 DAG,执行器跑 Python 插件 |
| 定时/周期任务 | 调度器 | Cron 表达式触发,调度器自动创建任务 |
| 脚本批量执行 | 执行器 | 通过 ScriptRunner 执行 Shell/Python 脚本 |
配置参考
# config.yaml
mode: executor # scheduler / executor
edsm:
transport_type: redis # redis / rabbitmq / inprocess
dsn: "redis://localhost:6379/0"
redis:
host: localhost
port: 6379
db: 1
mysql: # 调度器模式需要
host: localhost
port: 3306
user: root
password: ""
database: lzm_workers
executor:
max_concurrency: 10 # 最大并发任务数
heartbeat_interval: 3 # 心跳间隔(秒)
health_port: 9100 # 健康检查端口
scheduler:
reaper_interval: 30 # Reaper 巡检间隔(秒)
项目目录
lzm-workers/
├── pyproject.toml
├── src/lzm/workers/
│ ├── __init__.py # 导出 Scheduler / ExecutorRunner / Task
│ ├── cli.py # CLI 入口
│ ├── common/ # 公共模块
│ │ ├── config.py # 统一配置模型
│ │ ├── models.py # Task / TaskResult / TaskStatus
│ │ └── exceptions.py # 异常体系
│ ├── executor/ # 执行器板块
│ │ ├── runner.py # 主循环
│ │ ├── task_runner.py # 单任务引擎(租约+超时)
│ │ ├── pool.py # 执行器池(路由+并发)
│ │ ├── base.py # BaseExecutor 抽象基类
│ │ ├── plugin_executor.py # PluginExecutor → lzm-plugin
│ │ ├── script_runner.py # ScriptRunner → subprocess
│ │ ├── heartbeat.py # 心跳上报
│ │ ├── result_reporter.py # 批量回写
│ │ └── health_server.py # 健康检查 HTTP
│ └── scheduler/ # 调度器板块
│ ├── scheduler.py # Scheduler 入口
│ ├── task_manager.py # 任务 CRUD + 定时
│ ├── dispatcher.py # 执行器发现 + 分发
│ ├── reaper.py # 收割者
│ └── models.py # 数据库模型
├── tests/ # 26 个测试用例
└── docs/
└── logic-records/ # 架构设计文档
依赖
- lzm-edsm >= 0.2 — 事件驱动状态机引擎
- lzm-plugin >= 0.2 — 通用插件化框架
redis>= 5.0 — 心跳 + 租约 + 执行器发现
相关项目
| 项目 | 说明 |
|---|---|
| lzm-edsm | 事件驱动状态机引擎 |
| lzm-plugin | 通用插件化框架 |
| lzm-space | 空间存储编排引擎 |
| lzm-workers | 👈 本包:调度与执行计算框架 |
License
MIT © Lzm
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
lzm_workers-0.1.0.tar.gz
(30.0 kB
view details)
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 lzm_workers-0.1.0.tar.gz.
File metadata
- Download URL: lzm_workers-0.1.0.tar.gz
- Upload date:
- Size: 30.0 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.10
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
53877e31e44dde8cadbaacc82013bf41a05129ea3e6bb533b160da7d24c051ea
|
|
| MD5 |
7695c9d893ac9d4495619905e581dd4d
|
|
| BLAKE2b-256 |
dfd9bbabfa7ae1f3b37968d6dcae804e880f4675ed6482afe9ae5b9f1aa66422
|
File details
Details for the file lzm_workers-0.1.0-py3-none-any.whl.
File metadata
- Download URL: lzm_workers-0.1.0-py3-none-any.whl
- Upload date:
- Size: 37.9 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.10
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
3151786d11db98dcb3b9eae956aea7bd98845a984e88882cd5e64d67b1fe8de3
|
|
| MD5 |
b236c2b2c8bcbc7dc645b3d472e70665
|
|
| BLAKE2b-256 |
fc3b1cc6d32f93f10e6ae8ef0d4f003544840419a40f9b42096fa5b30ef3d650
|