P2 Foundry
管道构建:清洗与调度
把"取数 → 清洗 → 落表"写成一串有序 SQL 步骤,一次运行、多次重放,失败有 error、历史有记录;write_mode 决定清洗结果如何同步进目标对象表,cron 声明表达定时意图(当前需外部调度接入)。看完这 4 个故事,你就能自己把一张脏表清洗成可被对象消费的干净数据。
数据工程师
SQL 步骤
真实运行
定时调度
血缘联动
共 4 个故事
能 / 不能速览
✅ 这个主题能做
- 用有序 SQL 步骤(select / insert / update / create_table_as)声明式建管道
- POST /pipelines/:id/run 真实执行,连真实数据源顺序跑步骤,返回 rows_in / rows_out
- 运行历史(/pipelines/:id/runs)保留每次成功 / 失败记录与 error 详情,可排查重跑
- cron 声明 + enable / disable 表达"每天几点跑",停用后拒绝执行
- 运行成功自动打血缘:源表 → 管道,目标对象映射后管道 → 对象
⛔ 这个主题做不了
- 管道只支持 SQL 步骤,Python / 自定义代码步骤属 V2 未实现
- 跨源管道不支持:一条管道只能连一个数据源,步骤在同一源内执行
- cron 是"声明式"配置:平台不内置调度器,到点触发需外部调用 /run 接口
- 管道失败不自动回滚已写入的数据,也没有自动重试
适用角色
本主题面向两个角色:
- 数据工程师:建管道、写 SQL 步骤、排错重跑、配置调度——是核心使用者。
- 运营人员:点"运行"触发、看运行历史、确认管道健康度,不直接写 SQL。
下游的业务分析师消费管道产出的对象与指标;平台管理员在审计里能看到每次运行的触发人与结果。
能力速览(能做什么)
SQL 步骤管道
select 取数、insert / update 写入、create_table_as 落中间表,步骤有序执行,一条管道一个数据源。
真实执行
RunPipeline 经连接器连真实数据源顺序跑步骤,返回 status、rows_in、rows_out、error,不是纸面定义。
运行历史
每次运行都留记录(id 倒序),成功 / 失败、读入 / 写出行数、错误原因一目了然。
触发管理
manual / cron 两种触发类型,enable / disable 开关,停用的管道拒绝执行。
血缘联动
运行成功自动打点源表 → 管道边,目标对象映射后接上管道 → 对象边,数据来路可查。
调整指南(怎么调整)
- 改步骤:编辑管道时"sql_steps"可增删改,步骤名 / 类型 / SQL 均可调,保存即生效。
- 改调度:触发类型改 cron 时必须填 cron_expr(5 段标准格式),否则保存报错。
- 改写入:write_mode 支持 append(插入)/ append_only_new(按主键仅插新行)/ snapshot_replace(全量替换)/ snapshot_replace_and_remove(替换并移除已删行)/ changelog(P2 未实现),默认 snapshot_replace,创建时枚举校验;配了目标对象(target_object_type_id)且质量检查通过后,结果表会按 write_mode 同步进对象 base_table。
- 临时下线:点"停用"(disable),之后再点"运行"会返回"pipeline is disabled"。
- 接对象:先建好目标对象再填 target_object_type_id,血缘才能接上"管道 → 对象"。
做得好的场景
SQL 步骤管道把"清洗逻辑"变成可运行、可重放、可留痕的资产,特别适合以下场景:
- 固定清洗流程:过滤脏数据、字段转换、落中间表,一套 SQL 反复重放,不靠手工导表。
- 数据问题排查:失败时 error 精确到"第几步、什么 SQL、什么原因",运行历史按序可查。
- 例行调度声明:cron 表达式 + 启用状态把"每天 2 点跑"固化在定义里,交接一目了然。
- 接对象供业务:管道产物映射成对象,业务直接用对象查询,数据来路有血缘可追。
限制与不足
以下是明确的边界,使用前先知道:
- 仅 SQL 步骤:select / insert / update / create_table_as 四种,Python / 自定义代码步骤属 V2 未实现。
- 不支持跨源:一条管道只能连一个数据源,跨源聚合需拆成多条管道 + 中间表。
- cron 是声明式:平台保存并展示 cron_expr 与 enabled,但不内置调度器,到点触发需外部调度调 /run。
- 失败不回滚:步骤失败不会撤销前面已写入的数据,重跑前注意表状态;无自动重试。
- demo 数据重建:中间表建在演示 SQLite 里,重启后需重跑管道重建。
场景故事
故事 1
张工建第一条 SQL 管道:select 校验 + create_table_as 落表,真实运行
场景:管道搭建
角色:数据工程师
耗时:约 10 分钟
- 背景
- 张工要把"非取消订单"清洗成一张独立中间表 orders_active,供后续对象与指标消费。他在"管道构建"页新建一条 SQL 步骤管道,选演示数据源 foundry_demo_warehouse,声明两步:select 校验 + create_table_as 落表,然后真实运行。
- 传统做法对比
- 以前清洗要么 DBA 手动跑一段 SQL(无记录、无审计),要么写脚本定时跑(部署要半天);跑没跑、跑了几行全凭感觉。现在声明式管道保存即资产,点"运行"真实执行,rows_in / rows_out 当场可见。
- 角色
- 数据工程师(拥有管道构建与数据源权限);后续由业务分析师消费管道产出的对象。
- 操作步骤
-
- 登录平台(admin / admin1,端口 18081 或经 5173 的 /api 前缀),打开侧边栏"管道构建"
- 点"新建管道":名称=orders_active_sync、源连接器 ID 填 foundry_demo_warehouse 的 id(全新平台通常为 1)、写入模式=snapshot_replace、触发类型=manual、启用勾选
- 步骤 1:名称=check_valid_orders、类型=select,SQL 填 SELECT order_id, customer_id, amount, status FROM orders WHERE status != 'cancelled'
- 步骤 2:名称=create_orders_active、类型=create_table_as,SQL 填 CREATE TABLE orders_active AS SELECT order_id, customer_id, amount, status FROM orders WHERE status != 'cancelled'
- 保存后点"运行"
- 系统响应
- 运行返回示例:
{
"id": 3, "pipeline_id": 1, "status": "success",
"rows_in": 4, "rows_out": 4,
"started_at": "2026-08-09T10:15:00+08:00",
"finished_at": "2026-08-09T10:15:00+08:00", "error": ""
}
"运行历史"里新增一条 success;orders_active 表建成 4 行(5 条订单去掉 1 条 cancelled)。
- 结果洞察
- 声明式 SQL 管道把"清洗逻辑"变成可运行、可重跑、可留痕的资产:select 读到 4 行有效订单(rows_in=4),create_table_as 落表 4 行(rows_out=4),数字对得上。运行成功还自动打血缘"orders → orders_active_sync"(源表 → 管道)。以后任何时刻点"运行"即可重放整条清洗。
- 调整建议
- 步骤名起清楚(select / 落表)便于失败定位;SQL 别写死时间,用相对时间更通用;想接对象就在"本体工作台"先建好对象,再在管道里填目标对象类型 ID,结果表会按 write_mode 同步进对象表——重复运行看 write_mode 选对语义(append 会一直累加、snapshot_replace 每次全量替换)。
- 动手试一试
- 登录:admin / admin1。页面路径:管道构建 → 新建管道。输入内容:名称 orders_active_sync、源连接器 ID=1、步骤 1 select(WHERE status != 'cancelled')、步骤 2 create_table_as。预期结果:运行返回 status=success、rows_in=4、rows_out=4,运行历史出现一条成功记录。
- 限制提示
- 管道只支持 SQL 步骤(select/insert/update/create_table_as),Python 步骤属 V2 未实现;一条管道只能连一个数据源(跨源不支持);orders_active 建在演示 SQLite 里,demo 数据每次启动重建,重启后需重跑管道。
故事 2
管道跑挂了:字段名拼错,从运行历史定位并重跑
场景:失败排查
角色:数据工程师
耗时:约 8 分钟
- 背景
- 周三下午张工给 orders_active_sync 加了一列 payment_channel,想把支付渠道也洗出来,结果字段名拼错成 payment_channle,一运行就失败。他要定位错误、改定义、重跑,验证整条链路。
- 传统做法对比
- 以前脚本跑挂了只能翻日志、猜第几行,改错再跑一遍又要等调度;重跑前还得手工确认表状态。现在管道失败返回 error 精确到"第几步、什么 SQL、什么原因",运行历史保留每次记录,改一步、跑一步、看 error、再改,闭环很快。
- 角色
- 数据工程师(修改管道定义并重跑)。
- 操作步骤
-
- 在"管道构建"列表点 orders_active_sync 的"运行"
- 看返回 / 打开"运行历史",点开最新一条 status=failed 的 error 字段
- 确认是步骤 2 的 SQL 里列名写错
- 点"编辑",把 payment_channle 改回 payment_channel(并确认 orders 表确实有该列)
- 再次"运行"
- 系统响应
- 失败返回示例:
{
"status": "failed", "rows_in": 0, "rows_out": 0,
"error": "step 2 \"create_orders_active\" failed: no such column: payment_channle"
}
运行历史保留这条 failed 记录及其 error。
- 结果洞察
- 管道失败不是黑盒:error 字段精确指出"第 2 步 create_orders_active 失败,列 payment_channle 不存在",对照步骤列表一眼定位。改对字段名重跑,status=success。运行历史按 id 倒序保留每次运行(成功 / 失败都有 rows_in / rows_out / error),谁在什么时候跑挂过、挂在哪,全都有据可查。
- 调整建议
- 先跑只读的 select 步骤验证字段,再落 create_table_as;失败后用 GET /pipelines/:id 看当前 sql_steps,确认改到的是不是最新定义;目标表已存在时 create_table_as 会报错,先 DROP 或换表名。
- 动手试一试
- 登录:admin / admin1。页面路径:管道构建 → orders_active_sync → 编辑。输入内容:把某列名故意写错(如 amount → amunt)后运行。预期结果:run 返回 failed,error 为 no such column: amunt;改回正确列名重跑后 success。
- 限制提示
- 管道失败不自动回滚前面已写入的数据(无事务包裹),重跑前注意目标表状态;失败也不会自动重试;error 只记录本轮执行,运行历史里失败记录会一直保留,清理需在库侧处理。
故事 3
让管道"每天 2 点自动跑":cron 声明 + 启用状态
场景:定时调度
角色:数据工程师
耗时:约 5 分钟
- 背景
- orders_active 是每天都要用的中间表,张工不想总靠人点"运行"。他把 orders_active_sync 的触发类型改成 cron,表达式 0 2 * * *(每天凌晨 2 点),并保持启用,把"每天自动跑"固化进管道定义。
- 传统做法对比
- 以前定时任务要么配在服务器 crontab(谁配的、配没配都不知道),要么靠人每天手动跑(漏跑一天报表就错一天)。现在 cron 表达式 + 启用状态就写在管道定义里,交接的人一眼看懂"这条数据每天 2 点更新"。
- 角色
- 数据工程师(配置调度声明);运营人员依赖"每天都有新数据"。
- 操作步骤
-
- 在"管道构建"点 orders_active_sync 的"编辑"
- 触发类型改为 cron,cron_expr 填 0 2 * * *(5 段标准格式:分 时 日 月 周)
- 保存(cron 触发必须有 cron_expr,否则报校验错误)
- 确认列表行显示"cron · 0 2 * * *"与"启用"徽标
- 手动点一次"运行"验证 SQL 仍正常
- 系统响应
- 保存返回
{"code":0,"data":{"id":1,"enabled":true}};列表行显示 trigger_type=cron、cron_expr=0 2 * * *、状态=启用。
- 结果洞察
- cron 声明 + 启用把"每天自动跑"的意图固化在管道定义里,谁接管这条数据都看得懂。但要记住边界:当前版本调度是"声明式"——cron_expr 与 enabled 被记录和展示,实际到点触发需由外部调度(如 crontab / CI)调用 POST /pipelines/:id/run;平台自身不内置定时器。知道这个边界,就不至于误以为"到点自己跑"。
- 调整建议
- cron_expr 用 5 段标准格式:"每晚 2 点"写 0 2 * * *,"每小时"写 0 * * * *,"每周一早上 8 点"写 0 8 * * 1;临时下线用"停用"(disable),停用后 /run 会拒绝执行;外部调度器接入时把 triggered_by 带上,运行历史可追溯触发来源。
- 动手试一试
- 登录:admin / admin1。页面路径:管道构建 → 编辑 orders_active_sync。输入内容:触发类型=cron、cron_expr=0 2 * * *。预期结果:保存成功,列表显示 cron · 0 2 * * * 与"启用";不填表达式会报"cron trigger requires cron_expr"。
- 限制提示
- 平台当前只保存 cron 声明,不内置调度器,到点触发需外部调用 /run 接口(V2 会补内置调度);cron 声明本身不产生 run 记录,运行历史只在真正执行时生成;停用后 /run 返回 "pipeline is disabled"。
故事 4
管道结果被对象消费:源 → 管道 → 本体链路自动接上
场景:链路打通
角色:数据工程师
耗时:约 10 分钟
- 背景
- orders_active 表已经洗干净了,张工想让它变成业务对象,让分析师直接查。他在"本体工作台"建对象 order_active(主表 orders_active、主键 order_id),再把它设为管道的目标对象,跑一次管道,让"源 → 管道 → 本体"的血缘自动连起来。
- 传统做法对比
- 以前表和数据的关系全靠文档 + 人脑:"orders_active 是哪来的?"要翻聊天记录、问好几层人。现在对象 + 管道 + 血缘一条链,数据来路画成图,谁都能自己查。
- 角色
- 数据工程师(建对象 + 管道映射);业务分析师消费对象。
- 操作步骤
-
- 打开"本体工作台",新建对象 order_active:api_name=order_active、主表=orders_active、主键=order_id,映射 order_id / customer_id / amount / status 四个属性
- 编辑管道 orders_active_sync,目标对象类型 ID 填 order_active 对应的 id
- 保存并"运行"一次(成功)
- 打开"数据血缘"页:type=object、id=order_active、direction=upstream,查上游
- 系统响应
- 血缘返回示例:
{
"nodes": [
{ "type": "object", "id": "order_active" },
{ "type": "pipeline", "id": "orders_active_sync" },
{ "type": "source", "id": "orders" }
],
"edges": [
{ "from": { "type": "source", "id": "orders" },
"to": { "type": "pipeline", "id": "orders_active_sync" } },
{ "from": { "type": "pipeline", "id": "orders_active_sync" },
"to": { "type": "object", "id": "order_active" } }
]
}
对象查询选 order_active 返回 4 行有效订单。
- 结果洞察
- 数据流从此"看得见":orders 源表 → orders_active_sync 管道 → order_active 对象。张工把"表"变成"对象"后,业务同事不再问"orders_active 是哪来的"——血缘一查便知。之后若再给 order_active 建指标(对象 → 指标边在指标创建时自动打点),链路还会再往下延伸一级,形成源 → 管道 → 本体 → 指标的四级链路。
- 调整建议
- 目标对象映射填了之后,管道每跑成功都会再打点一次(幂等,不会重复插边);字段级映射(field_map)目前是表级 / 对象级为主,P2 会逐步细化;想展示"指标消费",给 order_active 建指标即可让链路延伸。
- 动手试一试
- 登录:admin / admin1。页面路径:本体工作台(建 order_active)→ 管道构建(填目标对象)→ 运行 → 数据血缘。输入内容:对象 order_active(主表 orders_active、主键 order_id)。预期结果:血缘页查 order_active 上游看到 source orders → pipeline orders_active_sync → object order_active 三级链路。
- 限制提示
- 结果表与对象 base_table 不一致时由 write_mode 目标对象同步落盘(append / append_only_new / snapshot_replace 等,changelog 标 P2 会报错);质量 error 级违规会跳过目标对象同步并把 run 标 failed(问题隔离);血缘边只在"运行成功"时打点,失败的 run 不产生血缘;demo 数据每次启动重建,重启后需重跑管道让血缘边重新补齐。
常见问题
管道到底"跑"了什么?是模拟还是真实执行?
真实执行。POST /pipelines/:id/run 会经连接器连到管道指定的数据源,按步骤顺序实际执行 SQL:select 用 ExecuteSQL,insert / update / create_table_as 用 ExecWrite,最后返回真实的行数与状态。
cron 配置了是不是到点就自动跑?
目前不是。平台把 cron_expr 与 enabled 作为"调度声明"保存并展示,但不内置调度器;到点触发需要外部调度(crontab / CI 等)调用 POST /pipelines/:id/run。V2 会补内置调度器。
管道跑失败了,已经写入的数据会不会回滚?
不会。管道步骤没有事务包裹,前面的步骤(如 insert)写入的数据会保留;失败时 run 标记 failed 并记录 error,需要你自己修正 SQL 后重跑。create_table_as 重建表时注意先 DROP 或换表名。
为什么管道只能连一个数据源?
当前管道是"单数据源内 SQL 步骤"(design 的 F3 决策,跨源不做)。跨源聚合需要拆成多条管道 + 中间表,分别清洗后再合并。
运行历史会保留多久?
运行记录持久保存在平台库(pipeline_runs),目前没有自动清理策略;列表按 id 倒序展示,最新的在最上面。演示环境数据量小,历史可长期保留供排查。
主题小结
一句话:管道构建器把"取数 → 清洗 → 落表 → 接对象"串成一条可运行、可重放、可追溯的 SQL 链路。记住几个边界:只支持 SQL 步骤、单数据源内执行、cron 是声明式(不内置调度器)、失败不回滚也不重试。