P2 Foundry
变更驱动流规则
让数据变更"主动"驱动业务动作:CDC 轮询源表发出 insert / update / delete 事件,流规则绑定 watcher、按条件表达式命中后自动执行 webhook / 通知 / 同步 / 质量画像 / 写回动作。看完这 3 个故事,你就能搭出"订单一进来就通知主管、大单就推外系统、脏数据自动修复"的自动化闭环。
运营人员
数据工程师
事件驱动
CDC
webhook
writeback
共 3 个故事
能 / 不能速览
✅ 这个主题能做
- 规则绑定 CDC watcher(cdc_watchers),订阅该 watcher 的全部变更事件
- 条件表达式过滤({{change_type}}、{{table}}、{{amount}}、{{pk_value}} 等变量),空条件恒真
- 五种动作:webhook(HTTP POST 10s 重试 1 次)/ notification(审计+站内通知)/ sync_run / quality_run / writeback
- writeback 走 writepath VALIDATE_AND_EXECUTE 五步流水线,不可绕过,幂等键 stream:<rule>:<event>
- 规则测试(POST /stream/rules/:id/test)与投递历史(fst_deliveries,同 rule 同 event 幂等去重)
⛔ 这个主题做不了
- CDC 是轮询比对式:检测延迟 = 轮询间隔(默认 30s,下限 5s),非实时日志流
- 进程内事件总线,跨产品不共享;重启期间的已落表事件不回放
- webhook 重试仅 1 次,非幂等场景可能重复送达,幂等由下游保证
- 投递历史无 status 过滤、不存事件体快照(只存 event_id / status / error)
- 重动作(webhook / writeback)在事件回调内同步执行,会阻塞 CDC 分发主流程
适用角色
本主题面向三个角色:
- 运营人员:配"订单插入 → 通知主管"这类规则,第一时间拿到数据变更信号,不用等日报。
- 数据工程师:配 sync_run / quality_run / writeback 等自动化动作,搭"事件 → 动作"编排。
- 平台管理员:维护 CDC watcher(监控哪张表、轮询间隔),通过投递历史排查失败动作。
规则执行的动作涉及平台其他服务(同步引擎、质量画像、writepath 写路径),权限与审计沿各服务自身语义生效。
能力速览(能做什么)
CDC watcher 绑定
规则关联 cdc_watchers:平台轮询源表(COUNT + max(watermark) 比对)发 insert / update / delete / snapshot 事件,规则订阅 watcher_id 即收事件。
条件表达式
复用 workflow 表达式语法:{{var}} 比较、数值运算、逻辑与或非、括号分组,payload 字段顶层展开,空条件恒真。
五种动作
webhook HTTP POST / notification 审计+站内通知 / sync_run 触发同步 / quality_run 触发质量画像 / writeback 走五步流水线写回。
幂等防重入
fst_deliveries 的 (rule_id, event_id) 复合唯一索引即幂等键:先 claim 后执行,同规则同事件重复触发被唯一冲突跳过。
测试与投递历史
POST /stream/rules/:id/test 构造模拟事件直接走处理流水线;GET /stream/rules/:id/deliveries 按时间倒序回查每次投递成败。
调整指南(怎么调整)
- 改触发范围:condition 用 {{change_type}} == 'insert' 限定事件类型、{{table}} == 'orders' 限定表、{{amount}} > 1000 限定数值,组合出精确触发面。
- 改动作参数:webhook 的 url / method / headers / body({{var}} 模板);writeback 的 action_name / object_type_id / params 指向本体 Action。
- 改启停:enabled=false 即退订,事件不再触发;删除规则为硬删(投递历史保留);外部直改库后可用 SyncSubscription 手动重载。
- 改投递预期:先点"测试"构造模拟事件验证命中与动作,再放心启用;测试会真实执行动作(webhook 会真的发请求),先确认目标地址。
- 改调度粒度:CDC watcher 的 interval_seconds 决定检测延迟(下限 5s),要更实时就调小间隔,注意别打爆源库。
做得好的场景
流规则让"数据变更"主动驱动业务动作,特别适合以下场景:
- 订单插入 → 通知主管:新订单落地即站内通知销售主管,替代等日报,抢时效。
- 大单 → webhook 推外系统:金额超阈值的订单实时推到 CRM / 消息系统,下游马上跟进。
- 脏数据自动修复回写:status='shipped' 但 ship_date 为空时自动补写,走 writepath 五步流水线安全回写。
- 主数据变更 → 下游同步:物料表变更自动触发 sync_run 刷新下游数据集;数据量大变更触发 quality_run 跑质量画像。
限制与不足
以下是明确的边界,使用前先知道:
- 检测延迟:CDC v1 是轮询比对(COUNT + max(watermark)),延迟 ≈ 轮询间隔(默认 30s、下限 5s),不是实时日志流。
- 重启不回放:订阅建立后只处理新事件,进程重启期间的 cdc_events 不补投。
- 同步分发:webhook 等重动作在回调内同步执行(最长约 20s),会阻塞 CDC 分发主流程。
- writeback 不可绕过:事件驱动写回与手工 Action 同等安全约束,无免校验快路径。
- 投递可观测有限:deliveries 无 status 过滤、不存事件体快照;事件详情需回 cdc_events 查。
- 注入依赖:sync_run / quality_run / writeback 的 runner / executor 未注入时动作记录 failed,规则仍可创建(延迟到执行时报错)。
场景故事
故事 1
订单插入 → 站内通知主管:不用等日报了
场景:事件通知
角色:运营人员
耗时:约 10 分钟
- 背景
- 运营希望订单表有新订单时第一时间通知销售主管,而不是等日报。小陈在 CDC 建了 watcher 监控 orders 表(水位列 updated_at),再到"流事件规则"建一条规则绑定该 watcher:条件 = 新行插入才触发,动作 = notification 站内通知,标题和正文里带上订单号与金额。
- 传统做法对比
- 以前新订单要靠日报或人工盯表才能知道,大客户首单可能要隔天才被发现。现在 CDC 轮询检出 insert 事件 → 规则条件命中 → 站内通知秒级到达主管,时效从"天"变成"分钟级"。
- 角色
- 运营人员(建规则 + 配通知);销售主管(收站内通知)。
- 操作步骤
-
- 确认 CDC watcher 已建(监控 orders 表)
- 侧边栏"流事件规则"新建规则"订单插入通知"
- watcher_id 选该 watcher,condition 填 {{change_type}} == 'insert' && {{table}} == 'orders'
- action_type=notification,action_config 填 {"channel":"internal","title":"新订单 {{pk_value}}","body":"{{table}} 新增订单,金额 {{amount}}","recipients":["u2"]}
- enabled 开启,保存后自动订阅该 watcher 事件
- 到 CDC 手动触发一轮 poll 验证
- 系统响应
- 创建返回
{"code":0,"data":{"id":"<uuid>","name":"订单插入通知","watcher_id":"aip:orders","condition":"{{change_type}} == 'insert' && {{table}} == 'orders'","action_type":"notification","enabled":true}};事件命中后投递记录落 fst_deliveries(status=success),规则 fire_count+1、last_fired_at 更新。
- 结果洞察
- 从"人查数据"变成"数据找人":订单表插入一行 → CDC 检出 insert 事件(pk_value=1001、payload 带 amount)→ 规则条件命中 → 审计 STREAM_NOTIFY + 站内通知落表 + 投递记录 success。销售主管不用等日报就能跟进新订单,通知标题里的 {{pk_value}} 被模板展开成真实订单号。
- 调整建议
- 条件里的 {{amount}} 是 payload 顶层字段,配 JSON payload 时注意键名一致;notification 渠道 foundry 侧未配飞书引擎时降级为"审计 + 投递记录"(不报错);先用"测试"按钮验证模板展开再启用,避免通知内容空变量。
- 动手试一试
- 登录:admin / admin1。页面路径:流事件规则 → 新建。输入内容:watcher 选 orders、condition={{change_type}} == 'insert' && {{table}} == 'orders'、action_type=notification、title="新订单 {{pk_value}}"。预期结果:点"测试"构造 insert 事件返回 matched=true、status=success;投递历史 +1。
- 限制提示
- CDC 检测延迟 = 轮询间隔(默认 30s),不是实时;条件未命中的事件不产生投递记录;测试会真实执行动作并落一条投递记录;事件回放未实现,重启期间的事件不补投。
故事 2
大额订单 webhook 推外系统:条件命中即触发,失败重试一次
场景:webhook
角色:数据工程师
耗时:约 12 分钟
- 背景
- 张工要搭"大单实时通知外部 CRM":订单金额超过 10000 的插入 / 更新事件,通过 webhook 推到外部系统。他建一条规则绑定 orders watcher,条件 {{amount}} > 10000 && ({{change_type}} == 'insert' || {{change_type}} == 'update'),动作配 webhook POST 到 CRM 接收地址。
- 传统做法对比
- 以前大单要人工盯着日报,或者写独立脚本定时扫表再调 CRM 接口,脚本挂了没人知道。现在事件驱动的 webhook 命中即推,投递历史里每次成败一目了然,目标 500 还会自动重试 1 次。
- 角色
- 数据工程师(配规则 + 排障投递);外部 CRM 系统(收 webhook)。
- 操作步骤
-
- 新建规则"大单推送 CRM",watcher 选 orders
- condition 填 {{amount}} > 10000 && ({{change_type}} == 'insert' || {{change_type}} == 'update')
- action_type=webhook,action_config 填 {"url":"http://crm:8080/hook","method":"POST","body":"{{change_type}} {{table}} {{pk_value}} {{amount}}"}
- 启用规则
- 在 CDC 制造一笔金额 15000 的订单,触发一轮 poll
- 到投递历史确认 status
- 系统响应
- 测试 POST /stream/rules/:id/test 返回
{"code":0,"data":{"rule_id":"<rule-id>","event_id":"<uuid>","matched":true,"status":"success","error":""}};投递历史返回:{"code":0,"data":{"deliveries":[
{"id":"<uuid>","rule_id":"<rule-id>","event_id":"<event-uuid>",
"status":"success","error":"","created_at":"..."}],
"total":1}}
目标 500 时重试 1 次后 failed,error="webhook 非 2xx 响应: 500"。
- 结果洞察
- webhook 动作配置即模板:body 里 {{change_type}} / {{table}} / {{pk_value}} / {{amount}} 运行时展开成真实事件字段。投递记录 fst_deliveries 是唯一"规则 → 事件"消费痕迹——同 rule 同 event 重复触发时 (rule_id, event_id) 唯一索引拦截,幂等跳过不重复推。
- 调整建议
- webhook 超时 10s、失败重试 1 次(共 2 次尝试),非 2xx 也算失败;目标地址要稳定可用,否则重试也失败会累积 failed 记录;非幂等 webhook 注意重试后的重复送达语义,幂等由下游自行保证。
- 动手试一试
- 登录:admin / admin1。页面路径:流事件规则 → 新建 → webhook。输入内容:condition={{amount}} > 10000,url=http://127.0.0.1:8080/hook(可本地起一个接收服务),body={{change_type}} {{table}}。预期结果:测试命中返回 matched=true、status=success,接收端收到真实请求体。
- 限制提示
- webhook 是同步执行,最坏情况阻塞 CDC 分发约 20s;重试仅 1 次不保证送达;投递历史不存事件体快照,事件详情回 cdc_events 查;条件表达式语法用 workflow/expression.go,写法与 Gateway 出边 Condition 一致。
故事 3
脏订单自动补 ship_date:writeback 走五步流水线安全回写
场景:自动写回
角色:数据工程师
耗时:约 15 分钟
- 背景
- CDC 检测到订单表有 status='shipped' 但 ship_date 为空的脏记录,需自动补 ship_date=now()。张工建流规则:条件命中脏订单事件,动作 writeback 指向本体 Action mark_shipped——执行走 writepath VALIDATE_AND_EXECUTE 全五步流水线,不可绕过。
- 传统做法对比
- 以前要么 DBA 手写 UPDATE 补数(无审计、无权限校验),要么等人发现脏数据再处理。现在事件驱动自动补数,且与手工 Action 同等安全:Schema 校验 → RBAC → 记录级权限 → 属性级写权限 → 强制审计 + 幂等登记,任何一步失败投递标 failed 并可回查。
- 角色
- 数据工程师(配规则 + 维护 Action);writepath 写路径(执行五步流水线)。
- 操作步骤
-
- 确认本体已有 Action mark_shipped(object_type_id=3)
- 新建规则"补 ship_date",watcher 选 orders
- condition 填 {{change_type}} == 'update' && {{status}} == 'shipped' && {{ship_date}} == ''
- action_type=writeback,action_config 填 {"action_name":"mark_shipped","object_type_id":3,"params":{"order_id":7},"user_id":"system"}
- 启用规则,触发一轮 poll
- 到投递历史确认写回成败
- 系统响应
- writeback 构造 writepath.ExecuteRequest(Mode=VALIDATE_AND_EXECUTE、IdempotencyKey="stream:<ruleID>:<eventID>")执行五步流水线;成功投递 status=success;权限不足返回
{"status":"failed","error":"写回未成功: failed (FORBIDDEN)"},投递历史保留失败记录供回查。
- 结果洞察
- 事件驱动的写回与手工 Action 走同一安全路径:参数 Schema 校验 → 逐动作 RBAC(action:mark_shipped,admin 绕过)→ 记录级权限 → 属性级写权限 → 强制审计 + 幂等登记 + 乐观锁回写 + 编辑态落盘 + 对象变化事件。任何一步失败事件都不会静默丢失——投递历史是回查入口,error_code 直接给出失败原因。
- 调整建议
- params 引用事件字段时用 {{var}} 模板(如 {{order_id}} 取 payload 里的行主键);user_id 决定写路径 RBAC 身份,权限不足先给该用户授权再重试;同 rule 同 event 幂等键保证重试不重复写;被 reject 的写回在投递历史里看 error_code 定位是哪一步拦截。
- 动手试一试
- 登录:admin / admin1。页面路径:流事件规则 → 新建 → writeback。输入内容:condition={{change_type}} == 'update' && {{status}} == 'shipped',action_config 指向 mark_shipped(object_type_id=3,params 带订单号)。预期结果:触发事件后投递历史 status=success,对象查询可见 ship_date 已补写。
- 限制提示
- writepath 执行器未注入时动作 failed("writepath 执行器未注入,writeback 动作不可用");五步流水线不可绕过,无免校验快路径;编辑态落盘后源表值不变(本体查询经编辑态叠加展示新值),需理解编辑态语义再决定是否依赖。
常见问题
事件是从哪里来的?
来自平台 CDC(platform/cdc):它按 interval_seconds 轮询源表(COUNT + max(watermark) 比对),发出 insert / update / delete / snapshot 四类事件。流规则绑定 cdc_watchers 的 watcher_id 订阅事件,进程内事件总线分发。
条件表达式支持哪些变量?
payload 字段顶层展开(如 {{amount}}、{{status}}),另有保留字段 {{change_type}} / {{table}} / {{watcher_id}} / {{watcher_name}} / {{datasource_id}} / {{pk_value}} / {{event_id}} / {{occurred_at}} / {{payload}}。空条件恒真。
同一条规则同一条事件会重复执行吗?
不会。fst_deliveries 的 (rule_id, event_id) 复合唯一索引是幂等键:先 claim 一条占位投递(status=failed、error=pending),同 rule 同 event 再次触发时插入唯一冲突被跳过(日志 debug"流事件幂等跳过")。
webhook 失败了会自动重发吗?
webhook 超时 10s、失败重试 1 次(共 2 次尝试),之后仍失败则投递记录 status=failed。没有异步重发/补偿机制;非幂等 webhook 注意重试可能重复送达,幂等由下游保证。
规则停用 / 删除后会怎样?
停用(enabled=false)即退订订阅,事件到达时 handler 重查规则发现已停用会忽略;删除是硬删规则(fst_rules),投递历史保留可查。规则增删改经 CRUD 自动同步订阅。
为什么说检测不是实时的?
CDC v1 是零外部依赖的轮询比对(COUNT + 水位),不是 Debezium 类日志流。检测延迟 ≈ 轮询间隔(默认 30s、下限 5s);且 insert 假定 pk 单调递增、delete 只发摘要不逐行 diff,原地 UPDATE 需配 watermark_col 才能检出。
主题小结
一句话:流规则在 CDC 事件总线之上提供"事件 → 规则 → 动作"自动化闭环——watcher 绑定 + 条件表达式过滤,命中后执行 webhook / notification / sync_run / quality_run / writeback 五种动作,幂等键防重入,投递历史可回查。记住几个边界:轮询检测有延迟、重启不回放、重动作同步阻塞、writeback 五步流水线不可绕过。