P1 AIP
工作流自动化:把查数变成自动跑的流水线
业务上很多事是"天天重复"的:早上看一眼销量、每周一汇总周报、库存低了提醒采购。工作流把这些动作编排成 DAG,一次搭好,手动 / 定时 / Webhook 三种方式反复触发。看完这 3 个故事,你就能把"查数→AI 总结→通知"跑成一条自动流水线。
运营经理
数据工程师
IT 管理员
DAG
cron
Webhook
共 3 个故事
能 / 不能速览
✅ 这个主题能做
- 8 类节点编排成 DAG:查询、NLQ、AI 分析、AI 文案、条件分支、通知、调外部 API
- 三种触发:手动运行、cron 定时、Webhook 公开端点(外部系统可直调)
- 节点间用 {{.node_id.output.field}} 模板传数据,AI 节点自动拿到上游查询结果
- 节点失败自动重试 2 次;条件分支未命中自动标记 skipped
- 执行全量落库:定义 / 执行 / 节点状态三张表,每次运行都能回看
- query_data 节点自动应用 RLS / CLS 安全策略,与聊天查数同一套规则
- 通知三渠道:internal(审计留痕)、feishu(飞书私聊)、email(SMTP 邮件真实发送)
⛔ 这个主题做不了
- 循环 / 人工审批 / 等待 / 并行分支 / Agent 任务五类高级节点未实现
- feishu / email 渠道需先配置好飞书凭证或 SMTP 参数,未配置而使用会让节点失败
- 每次运行是同步执行,没有异步任务队列和暂停 / 恢复
- cron 表达式非法时发布不报错(只 warn),该工作流只能手动 / 其他触发
- Webhook 是公开无鉴权端点,敏感操作不能只靠它
适用角色
本主题面向两个核心角色:
- 运营经理:把每天重复的"取数 + 汇总 + 同步"动作沉淀成工作流,是最大的受益者。
- 数据工程师:搭工作流定义、配 cron / Webhook、排执行失败——工作流的"搭积木"角色。
IT 管理员负责 admin 写操作权限(创建 / 发布 / 删除工作流);业务分析师直接消费工作流产物(报告、通知)。
能力速览(能做什么)
DAG 编排(/admin/workflows)
8 类节点连成有向无环图:trigger / query_data / nlq_query / ai_analysis / ai_generate / condition / send_notification / call_api。
三触发方式
手动 POST /workflows/:id/run;cron 5 段表达式定时;Webhook POST /webhooks/:name(body 作触发数据)。
模板数据流
{{.node_id.output.field}} 引用上游输出;input_mapping 给 AI 节点注入映射变量。
安全内置
query_data 节点 ValidateSQL 白名单 + RLS / CLS 注入,用户 ID 随执行上下文透传。
执行留痕
三表持久化:workflow_definitions / workflow_executions / workflow_node_executions,可取消执行。
通知渠道(send_notification)
channel=internal(审计留痕)/ feishu(飞书私聊,recipients 填 open_id)/ email(SMTP 邮件,recipients 填邮箱),主题与正文支持模板。
节点目录
GET /workflows/node-types 返回每类节点的配置字段 Schema,前端动态渲染表单。
调整指南(怎么调整)
- 改查询结果范围:query_data 节点 SQL 模板里加 LIMIT,或让 AI 节点把行数、金额作为输入去总结。
- 改触发频率:改 trigger_config 的 cron 表达式后需重新发布(发布时注册 / 反注册);重启服务会自动重载已发布的 cron。
- 改通知内容:send_notification 的 subject_template / body_template 引用上游输出,改模板字符串即可。
- 改分支条件:condition 节点支持 == / != / true / false,条件表达式里引用上游行数等字段。
- 改依赖关系:定义只能 draft 状态修改(先取消发布再改),改完版本自动 +1。
做得好的场景
工作流最适合"固定口径 + 反复执行"的报表动作:
- 日报周报自动化:查数 → AI 总结 → 通知一条链,每天 / 每周定时自动跑。
- 数据链路透明:每次执行每个节点的耗时、输出、重试次数全可回看,失败好排查。
- 与聊天同安全规则:工作流里跑的 SQL 同样过白名单 + RLS / CLS,不会绕过治理。
限制与不足
以下是明确的边界,使用前先知道:
- 高级节点缺失:无 loop / 人工审批 / 等待 / 并行分支 / Agent 任务,复杂编排做不了。
- 通知渠道依赖配置:channel=feishu / email 需先配好飞书凭证或 SMTP 参数才真实发送,否则节点直接失败;默认 internal 仅审计留痕。
- 同步执行:运行期间接口同步等待,长工作流(多个 AI 节点)响应会慢。
- Webhook 无鉴权:公开端点,需在网络层 / 网关限制调用来源。
- 无异步队列:不能排队执行、不能暂停 / 恢复。
场景故事
故事 1
运营经理把"每日销量查数 + AI 总结 + 通知"编排成工作流并手动跑通
场景:DAG 编排
角色:运营经理 + 数据工程师
耗时:约 10 分钟
- 背景
- 吴经理是运营负责人,每天上午都要让助理把演示库的订单数据拉出来、数一下总金额、再把结论发到群里,一套动作至少半小时。数据工程师陈工说:AIP 工作流能把这些串成一条流水线。两人打开 http://127.0.0.1:18080,用 admin / admin1 登录,进入"工作流编排"页(/admin/workflows)。
- 传统做法对比
- 以前每天人工查数 + 汇总 + 写总结 + 发消息,半小时到一小时,口径还得人记着;现在把三个节点连起来一次搭好,以后每次运行几秒钟出结果,结论还是 AI 自动写的。
- 角色
- 运营经理(业务视角,提出"要什么")+ 数据工程师(admin,负责创建工作流定义)。
- 操作步骤
-
- 进入 /admin/workflows 新建"每日销量日报"
- 加 query_data 节点 q:sql_template 填 "SELECT COUNT(*) AS cnt, SUM(sales_amount) AS total FROM orders"
- 加 ai_generate 节点 a:prompt 填 "用一句话总结:共 {{.q.output.rows.0.cnt}} 笔订单,总销售额 {{.q.output.rows.0.total}} 元"
- 加 send_notification 节点 n:body_template "{{.a.output.text}}"
- 连线 q → a → n,发布后点"运行"看结果
- 系统响应
- 运行返回执行结果(execution_id + 每节点状态 + 输出):
{
"execution_id": "a1b2c3d4-...",
"workflow_id": "wf-xxx",
"workflow_version": 1,
"trigger_type": "manual",
"status": "completed",
"duration_ms": 3120,
"nodes": [
{ "node_id": "q", "node_type": "query_data", "status": "completed", "duration_ms": 800 },
{ "node_id": "a", "node_type": "ai_generate", "status": "completed", "duration_ms": 1800 },
{ "node_id": "n", "node_type": "send_notification", "status": "completed", "duration_ms": 15 }
],
"outputs": {
"q": { "columns": ["cnt", "total"], "rows": [["5", "6470.5"]], "row_count": 1,
"sql": "SELECT COUNT(*) AS cnt, SUM(sales_amount) AS total FROM orders" },
"a": { "text": "今日共 5 笔订单,总销售额 6470.5 元。" },
"n": { "body": "今日共 5 笔订单,总销售额 6470.5 元。", "status": "sent" }
}
}
- 结果洞察
- 三个节点按拓扑序依次执行:query_data 查出 1 行(cnt=5、total=6470.5),ai_generate 用模板把行数拼进 prompt 让 LLM 生成一句话,send_notification 把结论写进审计 + 日志。吴经理发现每个节点的耗时都看得见(查询 0.8s、AI 1.8s),哪一步慢一目了然。
- 调整建议
- 模板里引用 rows 数组取字段要用 {{.q.output.rows.0.cnt}} 这种下标写法;想让总结更细,就在 ai_generate 的 prompt_template 里多引用几列;通知收件人列表用 recipients 数组配置。
- 动手试一试
- 登录:http://127.0.0.1:18080,admin / admin1。页面路径:/admin/workflows。操作:照步骤建"每日销量日报",发布后点运行。预期结果:返回 trigger_type=manual、status=completed,ai_generate 输出"今日共 5 笔订单,总销售额 6470.5 元"。
- 限制提示
- 通知默认走 internal 渠道只写审计与日志;想真实送达就把 channel 配成 feishu(飞书私聊)或 email(SMTP 邮件),并先在管理后台配好对应凭证,未配置而使用会让节点失败。SQL 模板里不要写 LIMIT 之外的复杂语句,执行最多返回 100 行;工作流每次运行是同步执行,AI 节点多了响应会变慢。
故事 2
设成 cron 定时任务:每周一早 8 点自动跑销量周报
场景:cron 定时
角色:数据工程师
耗时:约 5 分钟
- 背景
- "每日销量日报"跑通了,但吴经理不想每天早上手动点运行。陈工在触发器配置里加上 cron 表达式,让工作流每周一早上 8 点自动跑,跑完结论直接进审计,吴经理上班打开就能看。
- 传统做法对比
- 以前定时报表要么靠人记着每天手动跑,要么让 IT 在系统层面搭调度(排期至少一天);现在改一行 trigger_config 配置,发布后平台用 robfig/cron 自动注册,到点就跑。
- 角色
- 数据工程师(admin,改配置并发布)。
- 操作步骤
-
- 进入"每日销量日报"编辑页,trigger_config 填 {"type":"cron","cron":"0 8 * * 1"}
- 保存后点击发布(发布时校验定义 + 触发器,并注册 cron)
- 到 /workflows/:id/executions 查看执行记录,确认 cron 触发的执行已产生
- 系统响应
- cron 触发的执行返回 trigger_type=cron(触发时刻与表达式落在执行记录的 trigger_data_json 里):
{
"execution_id": "cron-run-001",
"workflow_id": "wf-xxx",
"workflow_version": 1,
"trigger_type": "cron",
"status": "completed",
"duration_ms": 2900,
"nodes": [
{ "node_id": "q", "node_type": "query_data", "status": "completed", "duration_ms": 700 },
{ "node_id": "a", "node_type": "ai_generate", "status": "completed", "duration_ms": 1700 },
{ "node_id": "n", "node_type": "send_notification", "status": "completed", "duration_ms": 12 }
]
}
- 结果洞察
- cron 触发走的是平台内置的 CronScheduler(robfig/cron/v3),带 SkipIfStillRunning 防重入——上次没跑完会跳过本次,不会两个实例打架。服务重启后 StartCron 会自动从数据库重载所有已发布且带 cron 的工作流,调度不丢。陈工在执行记录(workflow_executions 表 trigger_data_json 列)里核对了 triggered_at 时间戳,确认是定时触发的。
- 调整建议
- 想换频率改 cron 表达式后需重新发布(发布时先 Remove 旧的再 Add 新的,幂等);取消发布即可反注册定时;把全局变量(如 {{.global.report_name}})放进 global_vars 里统一管理。
- 动手试一试
- 登录:admin / admin1。页面路径:/admin/workflows → 编辑"每日销量日报"。输入内容:trigger_config 填 {"type":"cron","cron":"*/2 * * * *"}(每 2 分钟跑一次,便于观察)。预期结果:发布后 2 分钟内出现一条 trigger_type=cron 的执行记录;改回 "0 8 * * 1" 或取消发布可停止。
- 限制提示
- cron 表达式非法时发布只 warn 不报错(该工作流只能手动 / 其他触发);演示库每次启动重建,cron 定时跑的 SQL 查询的是重建后的演示数据;前端没有"立即测试 cron 表达式"的工具,配置前请自查语法。
故事 3
ERP 系统用 Webhook 触发"库存预警"工作流,命中条件才发通知
场景:Webhook + 条件分支
角色:数据工程师 + IT 管理员
耗时:约 10 分钟
- 背景
- 公司的 ERP 每次调价都会回写一份调价记录。陈工想:ERP 调完价,POST 一个 Webhook 到 AIP,触发工作流检查"调价后订单量是否异常",异常才发通知给采购。这样把外部系统直接接进了 AIP 流水线。
- 传统做法对比
- 以前外部系统对接要靠写接口、建消息队列、配置告警规则,一个完整闭环往往要 2~3 天;现在 ERP 只要 POST 一个公开地址,AIP 就把它当触发数据跑整条工作流,分钟级接入。
- 角色
- 数据工程师(搭 Webhook 工作流)+ IT 管理员(评估公开端点的网络限制)。
- 操作步骤
-
- 新建工作流"库存预警",trigger_config 填 {"type":"webhook","path":"price-alert"}
- query_data 节点 q:sql_template 填 "SELECT COUNT(*) AS cnt FROM orders WHERE product_id = {{.trigger.sku}}",data_source 填 aip_demo_warehouse
- condition 节点 c:condition_template 填 "{{.q.output.rows.0.cnt}} == 0"
- send_notification 节点 n:body_template "商品 {{.trigger.sku}} 近 30 天无订单,请确认是否停售"
- 连线 q → c → n(c 的 true_branch 指向 n),发布
- ERP 侧 POST /api/v1/webhooks/price-alert,body {"sku":"1"}
- 系统响应
- Webhook 触发的工作流执行结果(条件命中时):
{
"execution_id": "hook-run-007",
"workflow_id": "wf-hook",
"trigger_type": "webhook",
"status": "completed",
"duration_ms": 2100,
"nodes": [
{ "node_id": "q", "node_type": "query_data", "status": "completed", "duration_ms": 600 },
{ "node_id": "c", "node_type": "condition", "status": "completed", "duration_ms": 2 },
{ "node_id": "n", "node_type": "send_notification", "status": "completed", "duration_ms": 10 }
],
"outputs": {
"q": { "rows": [["0"]], "row_count": 1, "sql": "SELECT COUNT(*) AS cnt FROM orders WHERE product_id = '1'" },
"c": { "result": true, "expression": "0 == 0" },
"n": { "body": "商品 1 近 30 天无订单,请确认是否停售", "status": "sent" }
}
}
- 结果洞察
- Webhook 的 body 作为 trigger_data 注入,模板里用 {{.trigger.sku}} 取到 ERP 传来的商品号;condition 求值 "0 == 0" 为真,只放行了 true_branch 指向的通知节点——如果商品有订单,条件不命中,通知节点会标记 skipped,一条多余的告警都不发。没匹配到任何已发布 Webhook 工作流时返回 404。
- 调整建议
- 给不同业务配不同 path(如 price-alert / stock-alert),互不干扰;Webhook 无鉴权,建议在网关层限制来源 IP;如果查询或 AI 节点失败,可以再加一个 send_notification 节点接在失败路径上做异常通知。
- 动手试一试
- 登录:admin / admin1。页面路径:/admin/workflows 建"库存预警"并发布。操作:curl -X POST http://127.0.0.1:18080/api/v1/webhooks/price-alert -H "Content-Type: application/json" -d '{"sku":"1"}'(无需鉴权)。预期结果:返回 trigger_type=webhook 的执行结果,condition result=true,通知节点 completed;改用有订单的 sku(如 1 号商品)或改条件再试,观察 skipped。
- 限制提示
- Webhook 公开端点无鉴权,不能承载写入类操作;condition 只支持 == / != / true / false 四类求值;call_api 节点失败是非致命的(流程继续),如果你想让"调用失败也告警",得靠条件分支自己判断 status_code。
常见问题
工作流和 Agent / 工具系统是什么关系?
工作流是"人工编排的固定流水线"(DAG),按定义顺序执行;Agent 是按需调用工具做多轮推理。当前工作流引擎是独立模块,八类节点里 query_data / nlq_query 与工具系统共享同一套数据源与安全底座。
通知到底发到哪?
send_notification 支持三种渠道(channel 字段):internal(默认)把渲染好的主题与正文写进审计日志(事件类型 WORKFLOW_NOTIFICATION)并打印服务日志;feishu 按 recipients 的 open_id 推送到飞书私聊;email 按 recipients 的邮箱地址走 SMTP 发送邮件。feishu / email 需先在管理后台配好飞书凭证或 SMTP 参数,未配置而使用会让该节点直接失败(执行结果可见,不静默丢通知)。
能并行执行节点吗?
不能。MVP 按拓扑序顺序执行,多入边节点等全部上游完成;条件分支通过"边放行"实现,未命中分支的节点标记 skipped。并行分支(parallel 节点)未实现。
工作流里能改数据吗?
query_data 节点走 ValidateSQL 白名单,只放行 SELECT 类查询;写操作不是工作流的能力。要写数据请走数据源侧工具或后续的 Action 写路径。
取消执行会发生什么?
POST /workflows/executions/:eid/cancel 会把 pending / running 的执行置为 cancelled,未完成节点标记 skipped,数据库状态保持一致。已完成的执行不能取消。
主题小结
一句话:工作流 = 8 类节点连 DAG + 手动 / cron / Webhook 三触发 + 三表执行留痕。它把"查数 → AI 总结 → 通知"变成一次编排、反复自动跑。通知支持 internal(审计)/ feishu(飞书私聊)/ email(SMTP 邮件)三渠道。记住边界:无高级节点(循环 / 并行 / 人工审批)、同步执行、Webhook 无鉴权。SQL 安全(RLS / CLS)与聊天一致,不会被绕过。