干这行久了就会遇到一个很现实的问题:业务逻辑越来越复杂,手动跑脚本、用cron挂定时任务、靠人肉盯执行状态,这套传统玩法根本撑不住。你可能也遇到过——一个任务凌晨三点失败,第二天上班才发现;数据管道A依赖B,B依赖C,中间任何一环出错,整条链路塌方;任务重试、告警、日志追查全是手工作坊式的操作。我在接触deer-flow之前已经试过几个方案,要么太重,要么太绕,直到自己动手梳理这个项目,才发现工作流编排这件事,完全可以做得足够轻、足够直观。
deer-flow是一个面向中小型团队和个人开发者的分布式工作流编排引擎,核心解决的是"DAG流程调度、任务依赖管理、执行生命周期追踪"这三件事。它不需要像那些重型调度平台一样部署一整套复杂的微服务集群,而是用尽量收敛的组件完成流程定义、调度触发、任务分发和结果回收。本文会把整个项目的设计思路、核心模块、实际部署流程、常见坑点完整拆一遍,适合正在选型或打算自研工作流系统的读者参考。
1. 整体设计与思路拆解
1.1 这个项目到底解决什么问题
先说到底什么场景下需要deer-flow这样的东西。假设你维护一个数据报表系统,每天凌晨要从多个数据源拉取数据,清洗后写入数仓,再触发下游的报表计算和推送。听起来不复杂,但拆开看就是一条包含十几个任务的流水线,而且任务之间存在先后依赖。用cron写的话,只能做到"按时间触发",做不到"按依赖触发"。比如任务B要等任务A成功后才能跑,如果A跑了40分钟而B的cron设在A之后的固定时间点,那B要么空转,要么直接失败。
deer-flow解决的正是这个层次的问题:把任务定义成有向无环图(DAG)中的节点,节点之间的边表示依赖关系,调度器负责在满足依赖条件时自动触发下游任务。它不追求像大型数据平台那样管理几十万级别的任务量,而是把一个集群里几百个流程、几千个任务节点管得明明白白。对于一个日活几千到几万的业务系统来说,这个量级完全够用,而且运维成本低得多。
1.2 核心架构与模块划分
整个项目划分为五个模块,职责非常清晰。
- 调度引擎(Scheduler):负责扫描流程实例的触发条件,判断哪些任务节点可以进入待执行状态。调度引擎无状态,可以多实例部署,通过分布式锁保证同一时刻只有一个调度器在跑调度循环。
- 执行器(Worker):真正干活的节点,从队列里拉取待执行的任务,执行用户提交的脚本或命令,然后回调结果。Worker是无状态的,可以水平扩容。
- 控制台(Console):一个Web界面,用于流程定义、手动触发、查看执行日志和状态。面向使用方,不需要懂后端实现也能操作。
- 元数据存储(MySQL/PostgreSQL):保存流程定义、流程实例、任务实例、执行历史等关系型数据。这是整个系统的唯一事实来源。
- 队列与协调(Redis):承担两件事,一是任务队列的分发,二是分布式场景下的协调,比如调度器选主、节点心跳。
这个架构是不是有点眼熟?没错,它本质上是大厂调度系统的简化版:Airflow的DAG模型、DolphinScheduler的Master/Worker分离、加上轻量级的持久化和队列组件。区别在于,deer-flow刻意缩减了组件数量,没有引入ZooKeeper、没有独立的告警服务、没有单独的API网关,整个系统启动起来只需要三个服务加一个Web界面。
1.3 为什么是这套架构而不是别的方案
我在确定这个架构之前,其实做过几轮取舍。
第一轮考虑过完全复用Airflow,毕竟它的生态最成熟,DAG可以用Python写,插件也多。但Airflow的部署复杂度、Celery Executor相关的组件配置、还有它那套基于DAG文件解析的调度方式,对一个不太想投入太多运维精力的团队来说,重量感太明显了。
第二轮想法是自己从零写一套基于MongoDB和消息队列的系统,用JSON来描述工作流,这样可控性最强。但后来意识到,Redis和关系型数据库基本是每个后端团队都有的基础设施,没必要再引入新的中间件。而且MySQL存元数据、Redis做队列这个组合,在单集群几千任务节点的规模下,性能是完全够用的。
最终架构的核心思路是"保留两级抽象,去掉多余组件"。两级抽象指的是控制面(Console+Meta DB)和数据面(Worker),中间通过Redis队列解耦。这样的好处是:调度器和Worker可以独立扩缩容,控制面挂了不影响正在执行的任务,Worker挂了任务可以由其他Worker重新拉取。去掉的组件里,最典型的是不需要独立的告警模块,因为告警本质上就是"查询异常状态的任务实例",一个定时任务扫数据库再走一遍Webhook就够了。
2. 核心细节解析与实操要点
2.1 DAG流程建模与任务依赖
在deer-flow里,一个完整的工作流被称为流程定义(Flow Definition)。流程定义由若干任务节点(Task Node)和节点之间的依赖关系组成。依赖关系有两种表达方式。
第一种是显式依赖。节点B声明依赖节点A,那么调度器只有看到A的状态变成success后,才会把B置为ready状态。这种依赖用节点ID做引用,定义在JSON里就是一个edges数组。
{ "flowName": "daily_report", "nodes": [ {"id": "A", "type": "shell", "script": "pull_data.sh"}, {"id": "B", "type": "shell", "script": "clean_data.sh"}, {"id": "C", "type": "shell", "script": "compute_report.sh"} ], "edges": [ {"from": "A", "to": "B"}, {"from": "B", "to": "C"} ] }第二种是隐式依赖,靠的是数据可达性。比如任务B需要读取任务A产出的文件或表,开发者在定义B时声明数据来源,系统自动在A和B之间建立依赖。这种方式在数据管道场景里更直观,但实现起来需要额外维护一张数据血缘表。
我的建议是从显式依赖开始,尽量把依赖关系画清楚。原因很简单:工作流系统最难排查的问题就是"这个任务为什么跑了、为什么在这个时间跑了",显式依赖把原因写在了定义里,任何人打开DAG图都能看懂。隐式依赖到了一定规模后,排查依赖关系会变成一场噩梦。
2.2 任务节点的状态机设计
每个任务节点从被调度器感知到执行结束,会经历一个完整的状态机:
- init:流程实例刚创建,节点还没被调度器扫描。
- ready:所有上游依赖满足,节点进入待调度队列。
- running:Worker拉取了任务,正在执行。
- success:执行成功,结果已回收。
- failed:执行失败,进入重试判断逻辑。
- timeout:执行超时,系统主动中断。
- killed:被用户手动终止。
这个状态机是整个系统的核心,因为所有并发控制、失败重试、超时处理都围绕它展开。需要特别注意的一点是:状态迁移必须保证原子性。比如调度器扫描到某个ready任务,正要把它分发给Worker,此时用户手动终止了流程实例。如果没有锁保护,可能出现任务被分发出去了但流程实例已经被终止的脏状态。
实现上用Redis的分布式锁包住"状态检查+状态变更"的复合操作,锁粒度要控制在单个任务节点级别,不要锁整个流程实例。整条流程加锁会导致上下游节点全部串行化,在高吞吐场景下性能会很难看。
2.3 定时触发与事件触发
一个工作流系统不可能只支持手动触发,定时调度是刚需。deer-flow支持两种调度配置。
第一种是cron表达式,用标准的五位或六位表达式描述触发时刻。比如0 0 3 * * ?代表每天凌晨3点触发。cron表达式的解析不建议自己写,直接用现成的解析库(比如Go里的cron库或Java里的Quartz)就行,正则解析边界情况太多了,踩坑成本远大于引入一个成熟库的成本。
第二种是事件触发,监听消息队列或Webhook,收到特定事件后触发流程实例。我在项目里默认实现了基于Redis消息队列的事件监听,比如某个数据文件上传完成后向Redis发送一条消息,deer-flow收到后自动拉起对应的数据管道流程。
触发方式选哪个,取决于业务特性。定时任务适合周期性的、业务峰值不敏感的场景;事件驱动适合数据到达时间不确定、希望"到了就立刻跑"的场景。两个机制实现上有一定的重叠,但状态机层面完全一致——都只是生成一个新的流程实例,然后进入调度循环。
3. 实操过程与核心环节实现
3.1 环境准备与依赖部署
先把基础环境准备好。项目依赖三个外部组件:MySQL(或PostgreSQL)、Redis、以及一个可以跑HTTP服务的环境。如果是本地开发,直接用Docker起依赖最方便:
docker run -d --name deer-flow-mysql \ -e MYSQL_ROOT_PASSWORD=deerflow123 \ -e MYSQL_DATABASE=deer_flow \ -p 3306:3306 mysql:8.0 docker run -d --name deer-flow-redis \ -p 6379:6379 redis:7数据库起来之后,需要初始化表结构。项目里提供了一个init.sql脚本,包含流程定义表、流程实例表、任务实例表、调度日志表等。执行方式直接用mysql客户端导入:
mysql -h127.0.0.1 -uroot -pdeerflow123 deer_flow < docs/sql/init.sql初始化完成后,修改配置文件。配置项不多,主要包括数据库连接串、Redis地址、调度器扫描间隔、Worker并发数。贴一份核心配置示例:
server: port: 8080 database: host: 127.0.0.1 port: 3306 user: root password: deerflow123 name: deer_flow redis: addr: 127.0.0.1:6379 password: "" scheduler: scanIntervalSeconds: 5 # 调度器扫描周期 flowTimeoutMinutes: 120 # 流程实例超时时间(分钟) worker: concurrency: 10 # 单个Worker同时执行的最大任务数 heartbeatSeconds: 10 # Worker心跳间隔配置里的这两个参数要重点理解。调度器扫描间隔决定了"上游任务成功后,下游任务最快能多久被拉起"——扫描越频繁调度延迟越低,但会对数据库产生更多查询压力。我实测下来5秒的扫描间隔在几百个活动流程的规模下没有性能问题,可以当作默认值。Worker并发数决定了单个Worker能同时跑多少个任务,它受限于Worker机器本身的CPU和内存,需要根据实际任务负载来调。
3.2 通过控制台定义第一个工作流
依赖服务和配置都就绪后,启动调度器和Worker两个进程,然后打开控制台。控制台启动后,左侧是流程列表,右侧是画布区域。创建流程定义时,可以从左侧拖拽节点到画布上,也可以直接通过"代码模式"粘贴JSON定义。
我建议用代码模式来批量创建,特别是几十个节点的复杂流程,靠鼠标拖拽效率极低且容易连错边。控制台里的代码模式提供了JSON schema校验,格式有错误会直接标红提示。
创建一个最简单的三节点流程,目标是把数据库里的用户表数据每天凌晨导出为CSV文件,先做数据脱敏,再上传到对象存储。这个流程的JSON定义如下:
{ "flowName": "user_data_export", "description": "每日用户数据导出与脱敏上传", "scheduleCron": "0 0 3 * * ? *", "nodes": [ { "id": "export_db", "type": "shell", "script": "sh /opt/scripts/export_users.sh", "timeoutMinutes": 30, "retryTimes": 2, "retryIntervalSeconds": 60 }, { "id": "mask_data", "type": "shell", "script": "sh /opt/scripts/mask_users.sh", "timeoutMinutes": 20, "retryTimes": 3, "retryIntervalSeconds": 30 }, { "id": "upload_oss", "type": "shell", "script": "sh /opt/scripts/upload_users.sh", "timeoutMinutes": 10, "retryTimes": 1, "retryIntervalSeconds": 60 } ], "edges": [ {"from": "export_db", "to": "mask_data"}, {"from": "mask_data", "to": "upload_oss"} ] }这个定义里有几个参数需要专门讲一下。
timeoutMinutes是单个节点执行的超时上限。如果脚本卡住了(比如等待外部接口响应),超过这个时间后调度器会强制把节点置为timeout状态,然后把流程标记为失败。不设置超时的后果很严重——一个卡死的任务会一直占着Worker的并发额度,时间长了整个集群的并发能力会被拖垮。
retryTimes和retryIntervalSeconds是失败重试的配置。重试的默认策略是:任务失败后等待retryIntervalSeconds秒,然后重新进入ready队列。但要注意,重试是有限度的,重试次数用尽后节点才会进入failed状态。这里我踩过一个坑:给某个节点配置了5次重试,每次间隔60秒,一个本来应该快速失败的任务花了5分钟才真正失败。所以在配置重试时,要同时想清楚"这次失败是暂时的还是永久的"。连接超时、资源竞争这类失败适合重试;脚本语法错误、权限不足这类失败重试多少次都没意义。
3.3 工作流调度与执行的完整链路
流程定义保存后,就会进入调度器的管理范围。一个流程从被触发到执行完成的完整链路是这样的:
- 调度器根据cron表达式计算触发时刻。到点后,在流程实例表插入一条记录,状态为running,同时为流程定义里的每个节点创建对应的任务实例,状态为init。
- 调度器进入节点扫描循环。对于每个init状态的节点,检查其所有上游节点的状态。如果上游全是success,就把节点状态改为ready,并推入Redis任务队列。
- Worker通过拉模式从Redis队列获取任务。为什么要用拉模式而不是推模式?因为拉模式天然实现了负载均衡和背压控制——每个Worker根据自己的并发能力决定拉多少任务,处理完一个再拉下一个,不会出现推送模式下的队列堆积和Worker过载问题。
- Worker执行任务,把状态回写到数据库,标记为success或failed。如果是failed且还有重试次数,重置状态为ready并放回队列;如果重试次数用尽,节点置为failed。
- 当一个流程实例里所有节点都达到终态(success或failed)后,流程实例结束。存在failed节点时,整个流程实例标记为failed,触发告警逻辑。
这套链路的好处是每一步的职责都很单一。调度器只负责状态判断和入队,Worker只负责执行和回写,业务逻辑通过脚本和外部系统解耦。坏处是链路长,任何一个环节出问题都会导致任务卡住——这个问题在下一节详谈。
3.4 日志收集与执行追踪
任务跑起来之后,最重要的事情是能知道它执行得怎么样。deer-flow的日志分两个层次。
第一层是节点执行日志。Worker执行任务时会把标准输出和标准错误重定向到日志文件,文件路径规则是logs/{flowInstanceId}/{nodeInstanceId}.log。控制台的节点详情页会实时展示这个日志文件的内容,方便定位脚本执行问题。
第二层是调度日志。调度器每做一次状态判断、每次执行重试决策都会产生调度日志,存到数据库的scheduler_log表中。当出现"任务状态和实际不符"的问题时,查调度日志是唯一的突破口。我在排查问题时习惯先看调度日志,再看节点执行日志——先搞清楚"系统认为发生了什么",再看"实际发生了什么"。
还要强调一个开发体验上的细节:节点的script字段不要直接写复杂命令,写成sh /path/to/script.sh这种形式。原因有二,一个是JSON里嵌入多行命令需要大量转义,极其痛苦;另一个是脚本文件可以用版本管理来管,出了问题可以快速回滚。
4. 常见问题与排查技巧实录
4.1 任务一直处于ready状态但没被执行
这是使用过程中遇到最多的问题。现象是控制台上节点已经是ready了,但一直没有Worker拉取执行。排查思路按照下面顺序来。
第一步检查Worker状态。看控制台的Worker列表里有没有在线节点,如果Worker全部离线,检查Worker进程是否存活、Redis连接是否正常。很多时候是Redis密码配错了或者Redis内存打满了,Worker连不上Redis,心跳发不出去,被调度器判定为下线。
第二步检查任务队列。Worker在线但仍不拉取,大概率是队列数据出了问题。用Redis客户端直接查看队列长度:
redis-cli LLEN deer_flow:task_queue如果队列里有大量积压,说明Worker的消费速度跟不上生产速度。此时要么增加Worker实例,要么调大单个Worker的concurrency参数。如果队列是空的但控制台显示ready,说明状态同步出了问题。
第三步检查状态原子性。曾经遇到过一个情况:任务被调度器扫描到并发入队列,但Worker还没开始拉取时,用户通过控制台停止了流程实例。由于停止操作只改了节点状态为killed,队列里仍然残留了一条任务记录。Worker拉取后按节点当前状态判断,发现是killed就丢弃了,但这条消息没有走队列的ack机制,导致Redis里出现了一个永远无法被消费的堆积消息。解决方案是在Worker拉取后做一次二次校验,状态不是ready的节点直接丢弃结果,同时手动清理队列。
这个问题的本质原因是"数据库状态变更"和"Redis队列消息"之间没有做到事务一致。要根治需要在入队前加分布式锁,保证状态变更和入队操作原子执行;但考虑到实际场景中触发概率极低,通过二次校验也能达到同样的效果。
4.2 流程实例卡在running状态不结束
如果流程实例的running状态持续了几个小时,但所有节点都已经显示终态,说明流程实例的状态没有被正确汇总。调度器在节点状态更新后,会做一次流程级的状态汇总:过滤出所有节点的状态,如果全部success则汇总为success,如果有failed则汇总为failed。这个汇总操作由事件驱动,节点状态变更时会触发。
卡住的常见原因是某个节点的状态更新事件丢失了。比如Worker执行完任务,写数据库成功了,但在发送Redis事件通知时连接中断,导致汇总逻辑没有执行。这类问题很难在分布式系统里完全避免,所以我在系统里加了一个兜底机制:调度器每轮扫描时,除了扫描节点状态,也会找出"所有节点都是终态但流程实例仍为running"的记录,强制做一次汇总。
这个兜底扫描非常关键,它保证了即使事件通知链路有损耗,系统的最终状态也会收敛到正确值。建议任何做工作流系统的团队都把这条规则作为设计底线:对账逻辑一定要有,不能只依赖实时事件。
4.3 重试导致的重复执行问题
重试机制有个潜在的坑——任务可能被重复执行。举个例子:Worker执行脚本,脚本先写了数据库,准备返回结果时网络闪断,Worker和调度器的连接中断。此时Worker侧其实已经完成任务了,但调度器因为没收到结果,判定为失败,触发重试。重试后同样的脚本再跑一遍,数据就被重复写入。
这个问题被称为"at-least-once"语义,是分布式任务系统里最经典的问题。解决思路有两个层面。
第一个层面是业务幂等。脚本设计时自带幂等性,比如写入前先检查数据是否已存在,或者使用业务唯一键来去重。这个层面需要脚本开发者的自觉。
第二个层面是执行器层面的去重。在节点实例表里增加一个execution_token字段,每次重试前生成新token。如果脚本支持接收token并回传执行结果,调度器根据token判断是否为重复回执。这个实现起来有侵入性,但能覆盖更多场景。
我的实际建议是优先保证业务幂等,因为在跑批场景里,大部分操作本身就是"覆盖写"性质的,重复执行一次不会产生严重问题。只有对那种"累加型"操作,比如统计计数器累加、余额变更,才需要做执行器层面的强去重。
4.4 时间轮相关的问题
deer-flow的定时触发基于cron表达式,cron表达式的校验是一个很容易踩坑的点。比如很多人在测试时喜欢把表达式配成*/1 * * * * ?,也就是每秒触发一次。在系统刚部署时这样配不会有什么问题,但如果忘了改回去,系统会以每秒一个流程实例的速度创建数据,一晚上就能产生八万多个流程实例,数据库直接被打爆。
我在一个演示环境里遇到过类似情况,最后是清掉了几万条测试数据才恢复。所以建议在生产环境做一道保护:配置cron表达式时,如果同一流程在上一分钟内有未结束的实例,新的触发直接跳过。这相当于一个防重入锁,避免定时任务意外堆积。
另外时区问题值得提一下。数据库连接串和服务器时区如果不一致,cron计算的触发时刻会偏离预期。我习惯把所有服务器、数据库、Redis统一设为UTC+8,对业务方透明,排查问题时不用在脑子里来回切换时区。
5. 实操经验与调优建议
5.1 Worker数量与线程模型的平衡
很多初次使用的人会把Worker并发数调得很大,觉得并发越高越好。但并发数高并不总是好事。我遇到过一个场景:单个Worker并发调到30,结果任务脚本里有几个重的数据处理操作,直接把机器内存打满,触发了OOM,连带影响了其他任务的执行。
推荐的做法是先按单个任务的平均资源占用估算。比如一个任务平均占用500MB内存,机器总内存16GB,留出4GB给系统和Worker框架本身,那么并发数大概在20左右比较安全。这个数值需要跑几个任务后观察监控指标再微调,不要一开始就设很高。
另外,Worker的并发模型不要用简单的goroutine不限量启动,要做一个信号量控制的固定大小线程池。deer-flow里用的是channel加计数器的方式:每个任务先获取一个信号量slot,执行完再释放。这样即使队列里有大量任务等待,同时运行的也只有设置的那个上限。
5.2 MySQL连接池与查询性能
调度器每轮扫描会查询大量数据:待触发流程、待调度节点、过期实例。如果扫描间隔很短,数据库压力会集中在调度器这一侧。我在优化时做了两件事。
第一件事是加索引。任务实例表的查询模式主要是"按流程实例ID查节点"和"按状态查待调度节点",对应的索引分别是(flow_instance_id)和(status, flow_instance_id)。加上索引后,一个包含上千节点的流程实例状态汇总查询从原来的十几秒降到了几百毫秒。
第二件事是调整连接池参数。MySQL的max_connections要相应调大,同时调度器的数据库连接池要设置合理的最大连接数。如果调度器扫描频率是5秒一次,每次扫描需要查询上百条记录,连接池上限至少要覆盖这个并发查询量。我用的是连接池上限50,在日均上千个流程实例的场景下没有出现过连接等待超时。
5.3 流程定义的版本管理
生产环境的流程定义一定要支持版本管理。我在设计初期没有考虑这个问题,导致有一次修改流程定义时不小心把一个正在运行的节点的脚本路径改错了,新的执行直接报"文件不存在",老流程也受到牵连。
后来在流程定义表里增加了version字段,每次修改定义会生成一个新的版本,运行中的流程实例继续使用创建时的版本快照,新触发的实例使用最新版本。这样既保证了已运行任务的稳定性,也让变更可以回滚。这个功能实现起来并不复杂,核心是在创建流程实例时把整个流程定义的JSON快照存一份,而不是引用流程定义的当前状态。
5.4 告警策略的落地
最后说一下告警。我一开始想把告警做得很全面,比如节点失败、流程失败、调度超时、Worker失联全都告警。报警多了之后,团队会产生告警疲劳,真正重要的问题反而被淹没。后来收敛成三条规则:
- 流程实例最终状态为failed时告警(最多重试后仍失败)。
- Worker心跳丢失超过5分钟时告警。
- 任务节点在ready状态停留超过10分钟时告警。
第三条规则特别有用,它能在任务被"卡住"但还没有失败之前就发现问题,而不是等用户反馈才去排查。告警通道我直接接了团队的企业微信Webhook,消息内容包括流程名称、实例ID、失败节点、错误日志片段,收到告警的人可以快速判断问题严重程度,再决定要不要介入。
再分享一个细节:告警消息里一定要带流程实例的跳转链接,让人可以一键打开控制台查看详情。否则收到告警后还要去系统里翻找实例ID对应的流程,排查效率低不少。
6. 一些总结性体会
从需求梳理到部署上线,deer-flow这个项目做下来,我最大的体会是"工作流系统本身不难,难的是让它在不可靠的底层设施上保持可靠"。网络会断、机器会挂、脚本会卡死、人也会犯错,好设计不是假设一切顺利,而是假设每个环节都可能出问题,并有对应的兜底措施。
如果你正准备自建工作流引擎,我的建议是从最小闭环开始:先支持最简单的DAG调度和手动、定时两种触发方式,跑通之后再逐步加事件触发、版本管理、复杂告警。不要一上来就追求大而全——调度器、API、控制台、Worker、多租户、血缘、审计一口气全做,大概率半年内交付不了,而且你根本不知道哪块设计是真正符合自己业务需求的。
就写到这,希望能对正在做相关工作流编排系统的朋友有参考价值。在实际部署和使用过程中遇到的问题,欢迎交流。