Skip to main content

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 提供事件追踪、状态上报、持久化回放,并提供基于 Redis 的 demo 与 Go Worker 外部协作示例,用于展示按需构建跨进程、跨设备任务协作的方式。

项目结构(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;

    %% ===== Links =====
    TG --> CFB[CelestialFlow Web]
    CFB --> TG 

    style CFB fill:#ffeaf0,stroke:#d66b8c,stroke-width:2px,rx:10px,ry:10px

快速开始(Quick Start)

安装 CelestialFlow:

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

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

如果你只使用 CelestialFlow 的核心调度、可观测性与持久化能力,上面的安装已经足够。

如果你还需要启用 CelestialTree 事件追踪能力,则需要额外安装 celestialtree:

# 对已发布包使用者
uv pip install celestialtree

# 如果你是 clone 仓库后的开发者/贡献者
uv sync --group dev

一个简单的可运行代码:

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(name="DemoGraph")
    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;

    TM[TaskExecutor.md] --> TS[TaskStage.md] --> TG[TaskGraph.md]
    TM --> OB[BaseObserver.md]
    TM --> TME[TaskMetrics.md]

    TG --> TQ[TaskQueue.md]
    TG --> TN[TaskStages.md]
    TG --> TR[TaskReport.md]
    TG --> TSR[TaskStructure.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 execution;

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

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

📘案例教程

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

📚CelestialTreeClient

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

🎮demo/ 总览

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

🧪tests/ 总览

如果你想查看 bench 内容,这些数据也是框架中部分设计取舍的依据:

⚡bench/ 总览

环境要求(Requirements)

CelestialFlow 基于 Python 3.12+,默认运行时依赖以下核心组件。 其中 celestialtree 不再属于默认运行时依赖,而是额外安装的可选组件。

依赖包 说明
Python ≥ 3.12 运行环境,建议使用 3.12 及以上版本
requests HTTP 客户端库,用于任务状态上报与远程调用
tqdm 可选组件,进度条显示,用于任务执行可视化
  • 如需运行 demo/demo_redis.py 或 Go Worker 示例,请额外安装 redis 并准备 Redis 服务;这部分不属于默认运行时依赖。

  • 如需运行依赖 CelestialTree 的 demo / bench / 追踪查询,请额外安装 celestialtree,或直接在源码仓库中执行 uv sync --group dev。

  • 如需使用可视化的Web服务, 请额外安装 celestialflow-web 并运行 celestialflow-web --host 0.0.0.0 --port 5000。

文件结构(File Structure)

📁 CelestialFlow	(588MB 419KB 930B)
    📁 bench           	(296KB 806B)
        📁 [1项排除的目录]                  	(194KB 222B)
        🐍 bench_datastructures.py          	(6KB 690B)
        🐍 bench_execution_mode.py          	(2KB 707B)
        🐍 bench_futures_memory.py          	(2KB 269B)
        🐍 bench_gil_vs_nogil.py            	(10KB 121B)
        🐍 bench_graph_mode.py              	(7KB 450B)
        🐍 bench_hash.py                    	(7KB 67B)
        🐍 bench_hash_container.py          	(3KB 1009B)
        🐍 bench_hash_memory.py             	(3KB 642B)
        🐍 bench_http_grpc.py               	(2KB 521B)
        🐍 bench_ipc_queue.py               	(7KB 104B)
        🐍 bench_lock_overhead.py           	(9KB 421B)
        🐍 bench_mpqueue_vs_shared_memory.py	(13KB 127B)
        🐍 bench_observer.py                	(7KB 860B)
        🐍 bench_persistence_spout.py       	(4KB 340B)
        🐍 bench_queue.py                   	(5KB 857B)
        🐍 bench_requests.py                	(6KB 813B)
        🐍 bench_tqdm.py                    	(1KB 235B)
        🐍 bench_utils.py                   	(543B)
    📁 demo            	(145KB 1012B)
        📁 [1项排除的目录]  	(100KB 9B)
        🐍 demo_executor.py 	(1KB 495B)
        🐍 demo_funnel.py   	(2KB 289B)
        🐍 demo_graph.py    	(3KB 227B)
        🐍 demo_network.py  	(3KB 729B)
        🐍 demo_observer.py 	(4KB 270B)
        🐍 demo_redis.py    	(9KB 81B)
        🐍 demo_stages.py   	(4KB 219B)
        🐍 demo_structure.py	(11KB 593B)
        🐍 demo_utils.py    	(6KB 148B)
    📁 dist            	(302KB 861B)
        ❓ .gitignore                          	(1B)
        ❓ celestialflow-3.2.7-py3-none-any.whl	(79KB 235B)
        📦 celestialflow-3.2.7.tar.gz          	(67KB 836B)
        ❓ celestialflow-3.2.8-py3-none-any.whl	(83KB 244B)
        📦 celestialflow-3.2.8.tar.gz          	(72KB 569B)
    📁 docs            	(1MB 756KB 830B)
        📁 en[已折叠]   	(568KB 991B)
        📁 ja[已折叠]   	(642KB 383B)
        📁 zh-CN[已折叠]	(569KB 480B)
    📁 experiments     	(2KB 1021B)
        🐍 experiment_networkx.py	(1KB 884B)
        🐍 experiment_tqdm.py    	(1KB 137B)
    📁 img             	(5MB 871KB 242B)
        📷 file_structure.svg  	(4MB 918KB 1000B)
        📷 logo(old).png       	(836KB 542B)
        📷 logo.png            	(122KB 747B)
        📷 scc_condensation.svg	(17KB 1B)
    📁 src             	(1MB 873KB 44B)
        📁 celestialflow[已折叠]         	(1MB 852KB 902B)
        📁 celestialflow.egg-info[已折叠]	(20KB 166B)
    📁 tests           	(4MB 170KB 251B)
        📁 benchmark[已折叠]    	(39KB 852B)
        📁 funnel[已折叠]       	(96KB 363B)
        📁 graph[已折叠]        	(724KB 404B)
        📁 observability[已折叠]	(189KB 230B)
        📁 persistence[已折叠]  	(387KB 877B)
        📁 runtime[已折叠]      	(1MB 295KB 264B)
        📁 stage[已折叠]        	(381KB 203B)
        📁 utils[已折叠]        	(639KB 527B)
        📁 [1项排除的目录]      	(487KB 589B)
        🐍 conftest.py          	(1KB 38B)
        🐍 __init__.py          	(0B)
    📁 [12项排除的目录]	(573MB 949KB 913B)
    ❓ .env            	(468B)
    ❓ .gitignore      	(1KB 315B)
    📝 AGENTS.md       	(1KB 434B)
    ❓ LICENSE         	(1KB 65B)
    ❓ Makefile        	(155B)
    ⚙️ pyproject.toml  	(2KB 668B)
    📝 README.md       	(17KB 465B)
    🔒 uv.lock         	(121KB 572B)

celestial-flow 3.2.9

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

版本日志(Version Log)

  • 3.2.9
    • feat:
      • [IMPORTANT] 移除 stage 中 stage_mode,并添加 graph_mode
        • 破坏性更新
        • stage_mode 可以细粒度的控制每个 TaskStage 在图中的模式,但经过多年使用,我认为这种细粒度的控制并无必要,反而无谓的增加理解成本
        • graph_mode 则提供了更粗粒度的控制,用于统一控制所有 TaskStage 在图中是串行/多线程/并发运行,适用于大多数场景
        • 同时完善了 graph_mode(serial/thread/async) * execution_mode(serial/thread/async) 总共 9 种组合模式
      • 移除 schedule_mode
        • 破坏性更新
        • 这个模式带来了许多复杂度,但没有与原先的 stage_mode / execution_mode 产生明显的组合优势
      • 在 get_graph_analysis 中添加 graph_mode
        • 这个函数主要用于给web端提供信息
        • web端代码已经同步修改
      • 添加新的 warning 项,当图不是 dag 且 graph_mode=serial 时触发
      • 移除 task_executor 中的参数 persist_result
        • 这是为了简化代码逻辑,同时也使所有任务的最终状态(包括:成功/失败/未执行)都能被统一存储
      • 将 task.retry 事件降级,不再申请单独事件ID,不再在 lifecycle 中留痕
    • reafactor:
      • 合并 benchmark 相关函数对 sync 与 async 两种模式的执行
        • 性能无差异,代码好看一些
      • 将原本的 fallback 改名为 lifecycle
        • 早在重构原本的 fail 时就想起个更恰当的名字,这个版本才实现
      • 将 graph.source_lists 改为 graph.source_names
        • graph.source_lists 直接存储节点的引用,而 graph.source_names 存储节点的名称
      • 所有的 log / lifecycle 文件统一改名为 flow_log / flow_lifecycle
      • 删除 graph 中的 out_edges in_edges
        • 现在由 OrderGraph 负责管理边关系
      • 在 tarjan_scc 中使用 DFS 取代原先的递归逻辑
      • 将 util_graph 改名为 util_order_graph
    • fix:
      • 修复最后一次 report 无法正常上传的问题
      • 修复 graph.run / graph.run_async 中 is_put_signal 参数没有生效的问题
    • chore:
      • 更新文档

更多过往日志可看:

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

Release files for celestialflow 3.2.9

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for celestialflow 3.2.9
File Size Uploaded
celestialflow-3.2.9.tar.gz 72.2 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for celestialflow 3.2.9
File Interpreter ABI Platform
celestialflow-3.2.9-py3-none-any.whl Python 3 none any Details

Total release size: 156.3 kB

Release files / celestialflow-3.2.9.tar.gz

Download URL celestialflow-3.2.9.tar.gz
Size 72.2 kB
Tags Source
SHA-256 checksum
How to use checksums
6f76570923dc1154f295d701df08e5e53d8b86a09997198c3ff71fc45310762c
BLAKE2b-256 checksum
How to use checksums
efaf1c42c5acf624b76d07813d5269ea40cbe8bbab4bfc74672d0bae15f2b45c
Upload date
Uploaded using Trusted Publishing?
What is 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}

Release files / celestialflow-3.2.9-py3-none-any.whl

Download URL celestialflow-3.2.9-py3-none-any.whl
Size 84.1 kB
Tags Python 3
SHA-256 checksum
How to use checksums
77efc0eb3757966ddf4b6af4a8332923442a4fc75e37943b62b0e5d33d5f57ed
BLAKE2b-256 checksum
How to use checksums
eb73404eee420f87709756880ad0e91fc3c2d60653229a4d85c7aab0e8103ddc
Upload date
Uploaded using Trusted Publishing?
What is 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}

Release history Release notifications | RSS feed

3.3.1

2 release files

3.3.0

2 release files

This release

3.2.9 This release

2 release files

3.2.8

2 release files

3.2.7

2 release files

3.2.6

2 release files

3.2.5

2 release files

3.2.4

2 release files

3.2.3

2 release files

3.2.2

2 release files

3.2.1

2 release files

3.2.0

2 release files

3.1.9

2 release files

3.1.8

2 release files

3.1.7

2 release files

3.1.6

2 release files

3.1.5

2 release files

3.1.4

2 release files

3.1.3

2 release files

3.1.2

2 release files

3.1.1

2 release files

3.1.0

2 release files

3.0.9

2 release files

3.0.8

2 release files

3.0.7

2 release files

3.0.6

2 release files

3.0.5

2 release files

3.0.4

2 release files

3.0.3

2 release files

3.0.2

2 release files

3.0.1

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page