Skip to main content

统一的计算节点框架 — 调度器(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


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)

Uploaded Source

Built Distribution

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

lzm_workers-0.1.0-py3-none-any.whl (37.9 kB view details)

Uploaded Python 3

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

Hashes for lzm_workers-0.1.0.tar.gz
Algorithm Hash digest
SHA256 53877e31e44dde8cadbaacc82013bf41a05129ea3e6bb533b160da7d24c051ea
MD5 7695c9d893ac9d4495619905e581dd4d
BLAKE2b-256 dfd9bbabfa7ae1f3b37968d6dcae804e880f4675ed6482afe9ae5b9f1aa66422

See more details on using hashes here.

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

Hashes for lzm_workers-0.1.0-py3-none-any.whl
Algorithm Hash digest
SHA256 3151786d11db98dcb3b9eae956aea7bd98845a984e88882cd5e64d67b1fe8de3
MD5 b236c2b2c8bcbc7dc645b3d472e70665
BLAKE2b-256 fc3b1cc6d32f93f10e6ae8ef0d4f003544840419a40f9b42096fa5b30ef3d650

See more details on using hashes here.

Supported by

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