Skip to content

Latest commit

 

History

94 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Flowlet

轻量的工作流、配置、分发与并行执行基础库。

安装

pip install -e .

或使用 pixi

pixi install

核心特性

  • 统一执行抽象ExecutableUnit / Kernel / Workflow 提供一致的业务逻辑封装
  • 分支分发Dispatcher 基于条件函数实现输入路由
  • 配置与逻辑分离BaseConfig / BaseConfigContainer / BaseBranchConfig 支持 JSON / TOML / MsgPack
  • 并行透明ThreadParallelWorkflow(IO 密集型)、CoroutineParallelWorkflow(异步 IO)与 RayParallelWorkflow(CPU/GPU 密集型)统一接口
  • 跨进程进度ProgressManager + MPProgressProxy / RayProgressProxy,本地/多进程/Ray 三端汇聚
  • Edge 模式ThreadEdgeNode / AsyncioEdgeNode / RayEdgeNode 提供有状态、命令式、异步执行的节点抽象
  • FSM 控制层FSMEdgeNode 用有限状态机驱动 EdgeNode,支持事件、守卫、错误迁移和运行历史
  • 运行观测LogManager / StreamManager / TelemetryManager 汇聚日志、stdout/stderr 和资源/指标事件
  • 语义占位符Default / Placeholder / Emptyholder / Voidholder 精确表达空值语义

快速开始

from flowlet import Kernel, ThreadParallelWorkflow
from flowlet.config import BaseConfigContainer

class MyConfig(BaseConfigContainer):
    threshold: float = 0.5

class MyKernel(Kernel[MyConfig, list]):
    config: MyConfig = MyConfig()
    
    def __call__(self, data: list[float]) -> list[float]:
        return [x for x in data if x > self.config.threshold]

# 即时执行
result = MyKernel.run([0.1, 0.7, 0.3, 0.9], threshold=0.6)
# → [0.7, 0.9]

# 并行执行
with ThreadParallelWorkflow(max_concurrent_tasks=4) as wf:
    results = wf.map(lambda x: x * 2, range(10))
# → [0, 2, 4, 6, 8, 10, 12, 14, 16, 18]

模块导航

模块 说明 文档
flowlet.executable_unit 可执行单元(ExecutableUnit / Kernel / Workflow) docs/executable_unit.md
flowlet.dispatcher 条件分发器(Dispatcher) docs/dispatcher.md
flowlet.strategy 策略模式(BaseStrategy / @mount) docs/strategy.md
flowlet.config 配置系统(BaseConfig / BaseConfigContainer / BaseBranchConfig) docs/config.md
flowlet.parallel_unit 并行工作流(Thread / Coroutine / Ray / RayPoolCreator) docs/parallel.md
flowlet.edge Edge 模式(ThreadEdgeNode / AsyncioEdgeNode / RayEdgeNode) docs/edge.md
flowlet.fsm 有限状态机控制层(FSMEdgeNode / StateMachineSpec) docs/fsm.md
flowlet.base.progress 进度管理(ProgressManager / Proxy) docs/progress.md
flowlet.base.logs / stream / telemetry 日志、输出流和遥测汇聚 docs/observability.md
flowlet.base 基础工具(语义占位符 / 类型注解) docs/base.md
完整示例、最佳实践、注意事项 docs/examples.md

完整架构概览见 docs/index.md

About

No description, website, or topics provided.

Resources

Stars

2 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages