Skip to main content

A flexible GRAPH-based task orchestration framework.

Project description

CelestialFlow ——一个轻量级、可并行、基于图结构的 Python 任务调度框架

CelestialFlow Logo

中文 | English | 日本語

CelestialFlow 是一个轻量级但功能完全的任务流框架,适合需要 复杂依赖关系灵活执行模型跨设备运行实时可视化监控 的中/大型 Python 任务系统。

  • 相比 Airflow/Dagster 更轻、更快开始
  • 相比 multiprocessing/threading 更结构化,可直接表达 loop / complete graph 等复杂依赖模式

框架的基本单元为 TaskExecutor,可独立运行,并支持三种执行模式:

  • 线性(serial)
  • 多线程(thread)
  • 协程(async)

TaskExecutor 实现了对任务的结果缓存,任务去重,进度条显示,多执行模式比较等功能,单独使用也很好用。

但除去直接使用 TaskExecutor,更重要的是使用其子类TaskStage。TaskStage 可以互相连接,形成具有上游与下游依赖关系的任务图(TaskGraph)。下游 stage 会自动接收上游执行完成的结果作为输入,从而形成明确的数据流。

TaskStage 的任务执行模式同样包含三种,与TaskExecutor中一致。

在图级别上,每个 Stage 支持两种上下文模式:

  • 线性执行(serial layout):当前节点执行完毕再启动下一节点(下游节点可提前接收任务但不会立即执行)。
  • 线程执行(thread layout):当前节点在主进程的独立线程中启动,适合 I/O 密集型任务和不可 pickle 的函数(如 lambda)。

TaskGraph 能构建完整的 有向图结构(Directed Graph),不仅支持传统的有向无环图(DAG),也能灵活表达 树形(Tree)环形(loop) 乃至于 完全图(Complete Graph) 形式的任务依赖。

在执行与调度之外,CelestialFlow 进一步引入 CelestialTree(简称: ctree) 事件追踪系统,为每一个任务及其衍生行为(成功、失败、重试、拆分、路由等)记录明确的因果关系。借助 ctree,可以从任意一个初始任务出发,完整还原其在 TaskGraph 中的传播路径与执行轨迹,使任务系统可以进行完整的追溯、分析、解释

在此基础上,CelestialFlow 支持 Web 可视化监控,并可通过 Redis 实现跨进程、跨设备协作;同时引入基于 Go 的外部 worker(通过 Redis 通信),用于承载 CPU 密集型任务,弥补 Python 在该场景下的性能瓶颈。

项目结构(Project Structure)

flowchart LR

    %% ===== TaskGraph =====
    subgraph TG[TaskGraph]
        direction LR

        S1[TaskStage A]
        S2[TaskStage B]
        S3[TaskStage C]
        S4[TaskStage D]

        S1 --> S2 --> S3 --> S1
        S1 --> S4

    end

    %% 美化 TaskGraph 外框
    style TG fill:#e8f2ff,stroke:#6b93d6,stroke-width:2px,color:#0b1e3f,rx:10px,ry:10px

    %% 统一美化格式
    classDef blueNode fill:#ffffff,stroke:#6b93d6,rx:6px,ry:6px;

    %% 美化 TaskStages
    class S1,S2,S3,S4 blueNode;

    %% ===== WebUI =====
    subgraph W[WebUI]
        JS
        HTML
    end

    style W fill:#ffeaf0,stroke:#d66b8c,stroke-width:2px,rx:10px,ry:10px
    style JS fill:#ffffff,stroke:#d66b8c,rx:5px,ry:5px
    style HTML fill:#ffffff,stroke:#d66b8c,rx:5px,ry:5px

    R[TaskWeb]
    style R fill:#f0e9ff,stroke:#8a6bc9,stroke-width:2px,rx:8px,ry:8px

    %% ===== Links =====
    TG --> R 
    R --> TG 
    R --> W
    W --> R

快速开始(Quick Start)

安装 CelestialFlow:

# 推荐使用 `uv` 管理依赖与环境
uv pip install celestialflow

# 不过也可以直接使用 `pip`
pip install celestialflow

一个简单的可运行代码:

from celestialflow import TaskStage, TaskGraph

def add(x, y): 
    return x + y

def square(x): 
    return x ** 2

if __name__ == "__main__":
    # 定义两个任务节点
    stage1 = TaskStage(name="Adder", func=add, stage_mode="thread", execution_mode="thread", unpack_task_args=True)
    stage2 = TaskStage(name="Squarer", func=square, stage_mode="thread", execution_mode="thread")

    # 构建任务图结构
    graph = TaskGraph()
    graph.set_stages(stages=[stage1, stage2])
    graph.connect([stage1], [stage2])

    # 初始化任务并启动
    graph.start_graph({stage1.get_name(): [(1, 2), (3, 4), (5, 6)]})

注意不要在.ipynb中运行。

👉 想查看完整Quick Start,请见Quick Start

深入阅读(Further Reading)

若你想了解框架的整体结构与核心组件,下面的参考文档会对你有帮助:

推荐阅读顺序:

flowchart TD
    classDef core fill:#e6efff,stroke:#3b82f6,color:#1e3a8a;
    classDef runtime fill:#e9f8ef,stroke:#22c55e,color:#14532d;
    classDef structure fill:#fff6e6,stroke:#f59e0b,color:#78350f;
    classDef execution fill:#f3e8ff,stroke:#a855f7,color:#581c87;
    classDef web fill:#ffeaea,stroke:#ef4444,color:#7f1d1d;

    TM[TaskExecutor.md] --> TS[TaskStage.md] --> TG[TaskGraph.md]
    TM --> TP[TaskProgress.md]
    TM --> TME[TaskMetrics.md]

    TG --> TQ[TaskQueue.md]
    TG --> TN[TaskStages.md]
    TG --> TR[TaskReport.md]
    TG --> TSR[TaskStructure.md]

    TR --> TW[TaskWeb.md]
    TN --> GW[Go Worker.md]

    class TM,TS,TG core;
    class TP,TME runtime;
    class TSR structure;
    class TQ,TN,GW execution;
    class TR,TW web;

以下三篇可以作为补充阅读:

如果你更喜欢通过完整案例理解框架的运行方式,可以参考这篇利用 TaskGraph 从零开始构建项目的教程:

📘案例教程

如果你对3.0.7版本加入的ctree_client与其功能感兴趣, 可以看看这一篇:

📚CelestialTreeClient

你可以继续运行更多的演示代码,这里记录了各个演示文件与其中的演示函数说明:

🎮demo/

​如果你想运行测试代码, 可以先查看如下文档内容:

🧪tests/

如果你想查看bench内容, 这里的数据成为框架中部分设计的决策依据:

⚡bench/

环境要求(Requirements)

CelestialFlow 基于 Python 3.12+,并依赖以下核心组件。 请确保你的环境能够正常安装这些依赖(pip install celestialflow 会自动安装)。

依赖包 说明
Python ≥ 3.12 运行环境,建议使用 3.12 及以上版本
fastapi Web 服务接口框架(用于任务可视化与远程控制)
uvicorn FastAPI 的高性能 ASGI 服务器
requests HTTP 客户端库,用于任务状态上报与远程调用
networkx 任务图(TaskGraph)结构与依赖分析
jinja2 FastAPI 模板引擎,用于 Web 可视化界面渲染
tqdm 可选组件,进度条显示,用于任务执行可视化
redis 可选组件,用于分布式任务通信(TaskRedis* 系列模块)
celestialtree 可选组件,用于任务状态上报与远程调用(ctree_client

文件结构(File Structure)

FileStructure
celestial-flow 3.2.3

(该视图由我的另一个项目CelestialVault中inst_file.FileTree.print_tree()生成。转换为图片则借助Carbon。)

版本日志(Version Log)

  • 3.2.3
    • feat:
      • [IMPORTANT] 前端仪表盘设置面板中新添 节点等待使用全局估计 开关, 开启后节点卡片中原 等待 量将被替换为 全局等待, 并在卡片中显示 全局剩余时间
        • 这个功能非常有趣, 开启前后能看到数值大幅波动
      • [IMPORTANT] 重写任务注入页面, 现在实用性远远高于之前版本
        • 可以给每个节点单独配置注入任务列表, 在一起发送
      • [IMPORTANT] 移除 TaskExecutorget_argsprocess_result 两个方法, 以及 unpack_task_args 属性
        • 这是为引入泛型必要的修改
        • process_result 功能很简单, 对func输出的result进行再次处理, 但事实上在输入func前对齐进行包装也能达到一样的效果
        • get_args 功能更为复杂, 可以直接将前一节点提供的 result 转为当前节点所需的 task 类型, 非常灵活; 但问题在于太过灵活, 导致使用心智负担很大
        • unpack_task_args 就是 get_args 带来的一项心智负担, 默认 get_args 会把 task 进行 (task, ) 包裹后发送给func, 开启 unpack_task_args 后则直接发送
      • [IMPORTANT] graph 的init参数中添加 name, 以与 executor stage 一致
        • 破坏接口破坏性更新
        • 在日志的 start_graph end_graph 与web端的 graph_anaylysis 中都有显示
      • 彻底移除前后端通信中的 graph_summary , 原本残余的 全局剩余时间 现在拆为各个节点的 全局等待全局剩余时间
        • 这里所说的 全局 意为根据图论关系, 由上游剩余的任务数估算下游总共能获得多少任务
        • 例如: 图关系 A -> B, A已处理任务2, 未处理任务3, B已处理任务4, 未处理任务2. 这意味着A成功的2个任务为B带来总共了6个任务, 那么我们可以据此估计B总功能获取"3/2*6=9"个任务, 因此B的 等待 任务数量为2, 但 全局等待 任务数量为5
      • 前端中添加部分提示气泡, 鼠标放上去后可以介绍相关信息, 例如本次新加入的节点 全局等待 的含义
      • 移除前端中节点卡片的拖拽功能
        • 这个功能是在最早加入web页面时添加的, 当时感觉很帅, 但现在有点玩腻了
      • 前端错误日志页面添加 任务注入 按钮, 可以把选定任务直接添加到任务注入页面中节点的代注入任务列表中
      • TaskSplitter 中添加 split_item 方法, 可自由定义
        • 原本的 splitter 非常依赖于 get_args, 现在通过 split_item 方法稍微弥补其缺失的灵活性
    • refactor:
      • [IMPORTANT] 引入泛型, 同时强制性要求py版本>=3.12
        • 泛型的引入使得代码更加类型安全, 同时也提高了代码的可读性
        • 3.12版本对泛型的表述非常直观
      • 移除前端代码中所有对 localStorage 的使用
        • 在已经有config配置文件的情况下意义不大, 反而会带来困扰
      • 将任务队列的drain操作从graph层移至stage层
        • 之前不能这样做是因为节点的 stage_mode 可能为 process, 此时在主进程持有的节点非真实运行的节点
        • 算是3.2.0版本带来的持久影响之一
      • 优化 TaskMetric 中对锁的使用
      • 移除 TaskEnvelope 中的 change_id() 方法, 现在默认envelope不可变, 同时 emit_retry_envelope 中不再把原有的envelope的id修改后继续提交给worker, 而是直接使用新生成的envelope
      • 移除错误日志中的 error error_repr task_repr 字段
      • 修改前后端通信中任务注入数据的格式, 以方便同时提交多个节点的任务数据
      • 前端代码中开启strict检查
      • 将前端中 injection.css 文件拆分成多个文件
      • 修改 config.json 的数据格式, 现在按照生效区域进行分类
      • 收紧前端中的数据类型
      • 添加 ReportTaskGraph 类型, 专门用于给reporter做类型声明
    • fix:
      • i18n本地化在部分字段上失效的问题
      • 直接点击仪表盘中错误数字跳转到错误日志页后, 设置面板显示的还是仪表盘页面的设置
      • 前端中为各项仪表盘请求添加 RequestSeq, 以避免前后发送两个请求, 但因为延迟关系, 先发送的请求的返回覆盖掉后发送请求的返回
      • 修复前端中各项空字段需要在第一次refresh才显示的问题, 体感上会导致"加载"很慢
    • chore:
      • 将文档更新翻译为英/日两语
        • 这项操作太耗token了

更多过往日志可看:

change_log.md

Star 历史趋势(Star History)

如果对项目感兴趣的话,欢迎star。如果有问题或者建议的话, 欢迎提交Issues或者在Discussion中告诉我。

Star History Chart

许可(License)

This project is licensed under the MIT License - see the LICENSE file for details.

作者(Author)

Author: Mr-xiaotian Email: mingxiaomingtian@gmail.com Project Link: https://github.com/Mr-xiaotian/CelestialFlow

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

celestialflow-3.2.3.tar.gz (1.6 MB view details)

Uploaded Source

Built Distribution

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

celestialflow-3.2.3-py3-none-any.whl (1.6 MB view details)

Uploaded Python 3

File details

Details for the file celestialflow-3.2.3.tar.gz.

File metadata

  • Download URL: celestialflow-3.2.3.tar.gz
  • Upload date:
  • Size: 1.6 MB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.10.6 {"installer":{"name":"uv","version":"0.10.6","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":null,"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

File hashes

Hashes for celestialflow-3.2.3.tar.gz
Algorithm Hash digest
SHA256 57a9149d86b30b0771badaaaeb728e2b1ea1855ee3a32503b89414855b64a7b8
MD5 5835cc924cb0df7b22148464839ab3e8
BLAKE2b-256 2b41b5500bb7365bdb50b3279fa25cf91df617541d09916511b6d5df8aa22dd5

See more details on using hashes here.

File details

Details for the file celestialflow-3.2.3-py3-none-any.whl.

File metadata

  • Download URL: celestialflow-3.2.3-py3-none-any.whl
  • Upload date:
  • Size: 1.6 MB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.10.6 {"installer":{"name":"uv","version":"0.10.6","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":null,"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":null}

File hashes

Hashes for celestialflow-3.2.3-py3-none-any.whl
Algorithm Hash digest
SHA256 2dc1401113aff90ce769d7030d1cf43c13546ee980f36e58bf23422127b4ca86
MD5 447f6da7cba4ddbfa9b0d2c45bac8fa7
BLAKE2b-256 b60377a791149852c7a7f9abceea06217e2153633c4a6c53f5cb5afbc5647418

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