参考 Langflow 架构思想的自研轻量 LLM 工作流编排平台。 用手写实现 Langflow 的四大精髓——DAG 图模型、组件契约、执行引擎、流式输出——每一个模块都能讲出设计取舍。
- 后端:Python 3.10+ / FastAPI / SQLModel / httpx(进程内 async,无 Celery)
- 前端:Vue 3 / Vue Flow / Vite(Schema 驱动的可视化编辑器)
- 存储:SQLite(零配置开发)/ PostgreSQL(Docker 部署);缓存:进程内 LRU(接口对齐 Redis)
| 机制 | 一句话 | 解决了什么 |
|---|---|---|
| 端口(Handle)机制 | 连线绑定到具体端口而非节点对 | 分支出口、Prompt 动态变量、连线类型检查的共同基础 |
| 波次并行调度 | 按「输入就绪度」切批,asyncio.gather 并行 |
两个 1.1s 的 LLM 节点并行 → 墙钟 1.99x 加速 |
| SKIP 哨兵 + optional 端口 | 未走的分支输出哨兵值,下游据此跳过 | 分支跳过与多分支合并的直觉语义 |
| 错误传播与分支隔离 | 失败顶点下游整支 SKIPPED,旁支照常 | 单节点失败不炸掉整张图 |
| 事件总线 seq + 历史回放 | 自增序号 + 历史 buffer + 每订阅者独立队列 | 迟到订阅 / 断线重连不丢事件 |
| SSE 双层流式 | 顶点状态 + LLM token 按 vertex_id 多路复用 |
一条连接同时跑 N 路打字机效果 |
| 内容寻址缓存 | key = hash(组件类型 + 参数 + 输入) | 上游未变的下游秒回 |
flowchart TB
subgraph FE["前端 · Vue3 + Vue Flow"]
direction LR
A1["拖拽画布<br/>端口 Handle 连线"]
A2["Schema 驱动<br/>配置面板"]
A3["SSE 实时<br/>运行视图"]
end
subgraph API["API 层 · FastAPI"]
direction LR
B1["Flow CRUD<br/>版本化"]
B2["运行触发"]
B3["SSE 事件流"]
B4["组件注册表"]
end
subgraph ENGINE["执行引擎 · 核心"]
direction TB
C1["Coordinator 调度器"]
C2["Partitioner 分区器<br/>切就绪批次"]
C3["Executor 执行器<br/>asyncio.gather"]
C4["VerticesManager 状态机"]
C1 --> C2 --> C3 --> C4
end
subgraph BUS["事件总线"]
direction LR
E1["seq 自增"]
E2["history 回放"]
E3["每订阅者 Queue"]
end
subgraph COMP["组件系统"]
direction LR
D1["声明式 Spec"]
D2["注册表"]
D3["内置 8 组件"]
end
subgraph INFRA["基础设施"]
direction LR
F1["SQLite / PostgreSQL"]
F2["LRU → Redis"]
end
FE --> API
API --> ENGINE
ENGINE --> COMP
ENGINE -.发布事件.-> BUS
BUS -.SSE 推送.-> FE
ENGINE --> INFRA
flowchart LR
R["POST /flows/{id}/run"] --> V["加载 Graph<br/>绑定版本号"]
V --> VAL{"拓扑 + 端口校验"}
VAL -- 非法 --> X["422 拒绝"]
VAL -- 通过 --> P["Partitioner 切批<br/>ready / to_skip"]
P --> SKIP["标记 SKIPPED<br/>推事件 + 落库"]
P --> G["asyncio.gather<br/>并行执行 ready 批"]
SKIP --> CHK
G --> CHK{"全部到达终态?"}
CHK -- 否 --> P
CHK -- 是 --> F["汇总 run_finished"]
G -. 顶点事件 .-> BUS["事件总线"]
SKIP -. 跳过事件 .-> BUS
F -. 终态事件 .-> BUS
BUS -. SSE .-> UI["前端实时渲染"]
| Langflow 真实模块 | 本项目对应 | 说明 |
|---|---|---|
src/frontend(React + React Flow) |
frontend/(Vue3 + Vue Flow) |
等价替换 |
src/lfx/src/lfx/graph/ |
backend/app/graph/ |
Graph/Vertex/Edge + Kahn 拓扑排序 |
src/lfx/src/lfx/execution/(coordinator/executor/partitioner/registry) |
backend/app/engine/ |
四件套一一对应 |
src/lfx/src/lfx/components/ |
backend/app/components/ |
契约 + 注册表 + 内置组件 |
src/bundles/*(按厂商拆包) |
内置组件(不拆包) | Langflow 拆包为按需安装省体积;自研版先内置,演进路径见下 |
src/backend/base/langflow/services/(服务层) |
app/core/runtime.py |
轻量单例装配,不上 DI 框架 |
src/sdk(Python SDK) |
REST API | 简化,可选做 |
src/langflow-stepflow(新一代引擎) |
不实现 | 作为架构演进谈资 |
# 1. 后端(默认 SQLite,无 LLM Key 时自动 Mock 模式,离线可用)
cd backend
pip install -r requirements.txt
uvicorn app.main:app --reload --port 8000
# 2. 前端(vite 代理 /api → 8000)
cd frontend
npm install
npm run dev # http://localhost:5173
# 3.(可选)接入真实 LLM:DeepSeek / Ollama / vLLM
export FLOW_LLM_DEFAULT_BASE_URL=https://api.deepseek.com/v1
export FLOW_LLM_DEFAULT_API_KEY=sk-xxx
export FLOW_LLM_DEFAULT_MODEL=deepseek-chat
# 或不设全局,在画布 LLM 节点的面板里单独填 base_url / api_keyMock 模式:未配置任何 API Key 时,LLM 组件返回确定性 mock 回复并模拟 token 流—— 无网络、无 Key 也能完整演示「节点点亮 → 分支跳过 → 流式输出」全流程。
docker compose up --build
# 前端 http://localhost:8080 后端 http://localhost:8000/docscd examples
curl -X POST http://localhost:8000/api/flows -H "Content-Type: application/json" -d @01-basic-llm.json
curl -X POST http://localhost:8000/api/flows -H "Content-Type: application/json" -d @02-branch-merge.json
curl -X POST http://localhost:8000/api/flows -H "Content-Type: application/json" -d @03-rag-qna.json01基础链路:输入 → Prompt 模板 → LLM → 输出02条件分支:If-Else 走单边分支(另一边自动 SKIPPED)→ Merge 收集 → 输出03RAG 问答:向量检索 top-k → 拼接上下文 → LLM 引用回答
cd backend
pip install -r requirements-dev.txt
pytest tests -q # 26 个用例:拓扑/图校验/引擎/分支/缓存/API/SSEdocs/ 下有三个可直接运行的脚本,把引擎内部过程打印出来(需先配好后端依赖):
cd backend
# 端口机制:打印每个顶点的端口取值 + SKIP 的两条传播路径
FLOW_LLM_MOCK=1 PYTHONPATH="$(pwd)" python ../docs/demo_ports.py
# 波次调度 + 事件总线:逐轮打印 (ready, to_skip),实测 1.99x 加速
FLOW_LLM_MOCK=1 PYTHONPATH="$(pwd)" python ../docs/demo_waves.py
# 事件总线:seq 回放 / 迟到订阅 / 心跳 / token 多路复用 / 断线重连 / LRU 淘汰
FLOW_LLM_MOCK=1 PYTHONPATH="$(pwd)" python ../docs/demo_bus.py| 端点 | 方法 | 作用 |
|---|---|---|
/api/flows |
GET/POST | 流程列表 / 新建(graph 变化即新版本) |
/api/flows/{id} |
GET/PUT/DELETE | 详情(含最新 graph)/ 更新 / 删除 |
/api/flows/{id}/run |
POST | 触发运行;wait: true 同步返回,否则返回 run_id |
/api/runs/{run_id} |
GET | 运行结果 + 每个顶点的输出/状态/耗时/缓存命中 |
/api/runs/{run_id}/events |
GET (SSE) | 执行事件流(支持迟到订阅回放) |
/api/components |
GET | 组件注册表(前端侧栏动态渲染) |
/api/health |
GET | 健康检查 |
运行触发参数:
POST /api/flows/{id}/run
{
"inputs": {"n1": {"text": "运行时覆盖文本输入节点"}},
"use_cache": true,
"version": null,
"wait": false
}SSE 事件类型:run_started / vertex_started / vertex_finished(含输出与耗时)/
vertex_token(LLM 逐 token,按 vertex_id 多路复用)/ vertex_failed / vertex_skipped / run_finished。
连线绑定到具体输入/输出端口而非节点对。这是三种能力的共同基础:
- If-Else 的
true/false双出口分支; - Prompt 模板的
accepts_dynamic_inputs:连线接到哪个 handle,{{变量}}就取哪个值; - 连线时端口类型检查(上游输出类型必须兼容下游输入类型,
any万能)。
保存/运行前做 Kahn 拓扑排序(最小堆保证同层确定性):有环直接 422; 要做 Agent 循环需显式建模 loop 节点 + 最大迭代次数。
class ComponentSpec(BaseModel): # 组件的"身份证"
type: str # "llm" / "prompt_template" / ...
inputs: list[InputSpec] # 名称/类型/默认值/必填/控件类型/是否可连线/是否可选
outputs: list[OutputSpec]
accepts_dynamic_inputs: bool # Prompt 动态变量端口
timeout: int # 单顶点执行超时- 为什么声明式而不是函数签名? 前端配置面板要结构化描述(名称/类型/默认值/下拉选项), 函数签名做不到;同时驱动参数校验与连线类型检查。一份 Schema 三处消费。
- 一个 Input 既是面板参数又是连线端口(Langflow 的经典设计): 未连线时自动回落到面板值 / 默认值。
- 注册表模式:启动时注册,
/api/components暴露给前端 → 新组件不改前端代码。 - 演进路径:Langflow 的 bundles 拆包(
pip install lfx-bundle-openai)—— 核心只留契约,厂商组件按需安装。自研版先内置,规模上来再拆。
运行请求 → 加载 Graph(绑定 flow_version)→ 校验
→ loop:
Partitioner 扫描 pending 顶点 →
ready = 所有上游到达终态且必填输入无跳过信号
to_skip = 必填上游 failed/skipped,或输入端口为 SKIP
asyncio.gather 并行执行 ready 批(async 事件循环天然贴合 I/O 密集的 LLM 调用)
完成即解锁下一批,事件实时推 SSE
→ 全部到达终态 → run_finished
- 错误传播:顶点失败 → 必填下游整支 SKIPPED,旁支照常执行(分支隔离)。
- SKIP 哨兵:If-Else 未走的分支输出
SKIP,分区器据此跳过该分支下游。 - 可选端口(optional):Merge 声明
items为可选 → 某条上游分支被跳过时 丢弃该连线、其余输入照常合并(对齐「收集实际走到的分支」的直觉语义)。 - 多入聚合:多条连线打到同一输入端口 → 值聚合为 list(按源节点 id 排序保证确定性)。
- 超时与异常:
asyncio.wait_for按组件 spec.timeout 控制;异常在执行器层 转为 FAILED 事件,绝不中断同批其他顶点。
flowchart LR
P["Coordinator / Executor<br/>bus.publish"] --> W{{"seq += 1"}}
W --> H["history 列表<br/>回放源"]
W --> Q["各订阅者 Queue<br/>实时 fan-out"]
H --> S["sse_stream"]
Q --> S
S --> R1["① 先回放 history"]
R1 --> R2["② 再接实时事件"]
R2 --> R3["③ 空闲发 : ping"]
R3 --> SSE(["SSE 单条连接"])
SSE --> L1["第一层<br/>顶点状态事件"]
SSE --> L2["第二层<br/>vertex_token"]
L2 --> M["按 vertex_id 分路<br/>还原 N 路 LLM 流"]
publish无条件写 history +put_nowait推各队列 → 迟到订阅者也能补齐全部事件;sse_stream先把 history 塞进自己的队列再挂订阅集 → 消掉「回放期间新事件丢失」的竞态;- 每个订阅者独立
asyncio.Queue,慢消费者不拖慢他人;finally里discard防泄漏; - 心跳用 SSE 注释行
: ping,配X-Accel-Buffering: no穿透 Nginx 缓冲; - 事件语义按「状态覆盖」设计(started 清零 → token 累加 → finished 用最终结果覆盖), 因此全量重放后前端状态收敛,token 不会拼重。
- 版本化:flow 保存不覆盖,插入新
flow_versions;运行绑定具体版本号 →「正在跑的流程被编辑了也不受影响」。 - 内容寻址缓存:key = hash(组件类型 + 参数 + 输入)。上游不变 → 下游直接命中,
重跑流程秒回;
vertex_runs.cached可观测。内存 LRU 实现,接口 async 对齐 Redis,零改动替换。
Dify 用 Flask + Celery 把执行放到任务队列,换取可靠性但牺牲延迟; Langflow 用 FastAPI 原生 async 换取低延迟和流式自然性——我的自研版选后者, 因为 LLM 调用是 I/O 密集型,async 事件循环天然贴合;等规模上来再加队列,这是 YAGNI。
| 维度 | Dify | Langflow | 本项目 |
|---|---|---|---|
| 后端框架 | Flask + Celery | FastAPI | FastAPI |
| 执行模型 | 任务队列异步 | 进程内 async 协程 | 进程内 async |
| 工作流表示 | JSON DSL | JSON 图(React Flow) | JSON 图(Vue Flow) |
| 流式 | WebSocket/SSE | SSE | SSE(seq + 回放) |
| 组件扩展 | 插件机制 | 组件类继承 + bundles | 组件注册表 |
| 类别 | 组件 | 说明 |
|---|---|---|
| io | 文本输入 | 流程入口;运行时可用 inputs 按顶点覆盖 |
| io | 结果输出 | 可视化出口,透传展示 |
| prompt | Prompt 模板 | {{变量}} 插值,变量名 = 连线端口名(动态端口) |
| llm | LLM(OpenAI 兼容) | base_url 可配 → DeepSeek/Ollama/vLLM;token 级流式;无 Key 自动 Mock |
| rag | 向量检索 | embedding + 余弦 top-k;配置 API 走真实向量,否则本地哈希向量(离线可用) |
| tools | HTTP 请求 | 通用 API 调用;body 支持 {{input}} 引用上游 |
| control | 条件分支 | 12 种算子;未走分支输出 SKIP → 下游自动跳过 |
| control | 合并 | 可选端口多入聚合为 list |
flow/
├── backend/
│ ├── app/
│ │ ├── api/routes/ # flows.py / runs.py / components.py
│ │ ├── graph/ # vertex / edge / graph(校验) / topo(Kahn)
│ │ ├── engine/ # coordinator / partitioner / executor /
│ │ │ # vertices_manager / events / cache
│ │ ├── components/ # base(契约) / registry / llm / prompt /
│ │ │ # rag / tools / control / io
│ │ ├── models/ # db.py(4 张表) / repo.py
│ │ ├── core/ # config / runtime(装配)
│ │ └── main.py
│ ├── tests/ # 26 个用例
│ └── Dockerfile
├── frontend/ # Vue3 + Vue Flow
│ └── src/
│ ├── components/ # FlowCanvas / FlowNode / Palette /
│ │ # InspectorPanel(Schema 驱动) / RunPanel(SSE)
│ ├── api.js / store.js
│ └── App.vue
├── examples/ # 3 个可导入的示例流程
├── docs/ # 3 个运行时演示脚本
├── docker-compose.yml
└── README.md
- 缓存换 Redis:
VertexOutputCache接口已 async 化,实现get/set即可全局共享; - 事件总线换 Redis pub/sub:多副本部署时 SSE 不再受单进程内存限制;
顺带把 seq 写进 SSE 的
id:字段,用Last-Event-ID做增量重放替代全量重放; - 组件拆包(对标 Langflow bundles):核心留契约,厂商组件独立 pip 包按需安装;
- 执行队列(对标 Dify Celery):重负载/多租户时把运行移出 Web 进程;
- Agent 循环节点:显式 loop 节点 + 最大迭代次数(拓扑校验白名单);
- 用户体系 / RBAC / 运行审批(Langflow services/ 的对应物)。