1. 页面概览

1.1 是什么

管道构建(PipelinesPage.vue,路由 /foundry/pipelines)是 LightFoundry 的「SQL 步骤声明式管道」管理界面,对标 Palantir Foundry 的 Pipeline Builder(Transforms)。与可视化 DAG 画布不同,本实现采用单数据源内有序 SQL 步骤:一条管道 = 一个源连接器(source_connector_id)+ 一串顺序执行的 SQL 步骤(sql_steps,类型 select/insert/update/create_table_as)+ 可选的写入模式(write_mode)+ 可选的触发方式(manual/cron)+ 可选的启用开关。管道可真实运行:后端建运行记录 → 连接真实数据源 → 顺序执行 SQL 步骤 → 质量规则检查 → 按 write_mode 同步到目标对象 → 落库运行状态与行数。

页面同时提供运行历史查看:选中管道可展开该管道的运行记录表(状态/rows_in/rows_out/错误),用于验收每次运行的结果。V5 Stage 0 起,cron 触发的管道在保存时同步注册到平台统一调度器(platform/scheduler),实现真正的定时调度。

1.2 核心价值

维度说明
声明式编排管道定义存 SQL 步骤 JSON,创建/编辑所见即所得,保存即生效
真实执行「运行」按钮触发后端真实执行:连数据源 → 顺序跑步骤 → 质量检查 → 按 write_mode 落目标对象
全链路留痕每次运行产生 run 记录(status/rows_in/rows_out/error/triggered_by),并可展开查看
写模式语义append / append_only_new / snapshot_replace / snapshot_replace_and_remove / changelog 五种写模式枚举,控制目标对象同步方式(changelog 未实现:创建/更新即被后端显式 400 拒绝,前端下拉禁用并标注「暂未实现」)
定时调度cron 触发 + 启用开关,配合平台统一调度器实现定时任务(V5 Stage 0)
质量联动管道 run 自动执行挂在该管道上的质量规则;severity=error 违规会使 run 标 failed 并阻断目标对象同步

1.3 一句话总结

在 Foundry 里用「数据源 + 有序 SQL 步骤 + 写模式」声明一条数据处理管道,页面负责定义/启停/手动运行/查看历史,后端负责真实执行与质量/血缘联动。

2. 访问入口

2.1 路由与菜单

2.2 认证与权限

2.3 端口与 API 前缀

3. 界面布局

┌──────────────────────────────────────────────────────────────┐
│ ① 页头:「管道构建」标题 + 「新建管道」/「收起表单」按钮         │
│ ② 操作结果提示条(alert 成功/失败,可「关闭」)                 │
│ ③ 创建/编辑表单卡片(点「新建管道」或「编辑」才展开)           │
│    ├─ 名称/源连接器 ID/目标对象类型 ID/写入模式/触发类型/cron  │
│    ├─ 描述 + 启用勾选                                          │
│    └─ SQL 步骤动态列表(+ 步骤 / 每步名称+类型+SQL+删除)       │
│       └─ [保存] [取消]                                         │
│ ④ 管道列表卡片(表格)                                         │
│    └─ ID/名称/描述/写入模式/触发/启用/操作(运行/继续轮询*/    │
│       停用/编辑/删除/运行历史)  *轮询超时后条件出现            │
│ ⑤ 运行历史卡片(选中管道展开)                                 │
│    └─ ID/状态/触发/rows_in/rows_out/开始/结束/错误             │
└──────────────────────────────────────────────────────────────┘

各板块职责:

4. 交互元素详解

4.1 页头与表单

位置元素含义默认值/必填操作效果触发的后端调用
页头「新建管道」/「收起表单」展开/收起空表单openCreate():清空 editingPipelineform,切换 showForm
表单名称(api_name)*管道 api_name必填,占位 如 orders_sync;编辑时 disabled文本输入
表单源连接器 ID *管道数据源 id必填(number),占位 数据源 id(如 1)数字输入
表单目标对象类型 ID可选,写入目标对象可空(number)数字输入
表单写入模式目标对象写模式默认 append;选项 append/append_only_new/snapshot_replace/snapshot_replace_and_remove/changelog(changelog 下拉禁用并标注「暂未实现」,创建/更新即 400)下拉
表单触发类型manual/cron默认 manual下拉
表单cron 表达式 *定时表达式trigger_type=cron 时显示且必填,占位 如 0 2 * * *文本输入
表单描述管道描述可空文本域
表单启用enabled 勾选默认勾选复选框
表单SQL 步骤(至少 1 个)+「+ 步骤」步骤动态列表至少 1 个有效步骤addStep() 追加空步骤
步骤行步骤名输入步骤 name必填(保存时校验)文本输入
步骤行变换类型下拉transform_type默认 select;选项 select/insert/update/create_table_as下拉
步骤行SQL 编辑框步骤 SQL必填(保存时校验),等宽字体、spellcheck=false文本域
步骤行「删除」按钮移除该步骤removeStep(i)
表单「保存」/「保存中...」提交表单busy 时禁用校验后创建或更新,成功关闭表单并刷新列表POST /pipelinesPUT /pipelines/:id
表单「取消」放弃编辑closeForm()

4.2 管道列表操作

位置元素含义操作效果触发的后端调用
列表行「运行」按钮手动触发管道confirm('确定运行管道"xx"吗?') 通过后执行:快路径 200 提示 管道运行已启动 run_id=<id> status=<status>;202 异步受理则提示 管道已受理,正在执行(task_id=…)退避轮询任务直至 success/failed/cancelled(成功后提示 管道运行成功 run_id=…,失败/取消 alert-error 展示原因),轮询期按钮显示「运行中」并禁用;若当前已展开该管道历史则终态后自动刷新POST /pipelines/:id/run(body {triggered_by:'web'})+ GET /system/tasks/:id
列表行「继续轮询」按钮(条件出现)自动轮询达 5 分钟上限仍未终态时,在该管道行露出,供用户续轮询点击复用上次 task_id/管道 id 重新开始退避轮询直至终态(resumeRunPolling);非该状态时不渲染GET /system/tasks/:id
列表行「停用」/「启用」按钮切换 enabled成功提示「管道已停用/已启用」并刷新列表POST /pipelines/:id/disablePOST /pipelines/:id/enable
列表行「编辑」按钮打开编辑表单openEdit(p):回填表单(含 sql_steps 深拷贝),showForm=true
列表行「删除」按钮删除管道定义confirm('确定删除管道"xx"吗?') 通过后执行;若该管道历史已展开则收起;刷新列表DELETE /pipelines/:id
列表行「运行历史」/「收起历史」展开/收起该管道运行记录首次点按加载;再次点收起并清空 runsGET /pipelines/:id/runs
列表列启用徽标status-enabled 启用 / status-disabled 停用只读

4.3 运行历史表

说明
ID运行记录 id
状态徽标 status-running/status-success/status-failed
触发triggered_by,页面触发为 web,cron 为 cron,缺省显示 -
rows_in读入行数(?? '-' 处理 0)
rows_out写出行数(?? '-' 处理 0)
开始时间 / 结束时间formatTime 格式化:2026-08-09T10:00:00Z2026-08-09 10:00:00(replace T + slice 19)
错误失败原因(红色 run-error),无错误显示 -

5. 后端关联

5.1 API 客户端

action/web/src/api/client.js:baseURL /api/v1、timeout 30000、Bearer aip_token、401 跳登录(同指标页,复用同一实例,PipelinesPage 直接 import apiClient from '../api/client.js')。

5.2 端点表

方法路径请求体要点超时
GET/pipelines30s
POST/pipelines{name, description, source_connector_id, target_object_type_id, sql_steps[], write_mode, trigger_type, cron_expr, enabled}30s
GET/pipelines/:id30s
PUT/pipelines/:id同创建 body 不含 name(name 不可改)30s
DELETE/pipelines/:id无(运行历史保留)30s
POST/pipelines/:id/run{triggered_by}(可空,空 body 用当前用户 id;页面传 web15s 内完成返回 200(run 信息);超时 202 {task_id,status,type}(任务化受理,非错误)
GET/system/tasks/:id无(后端 task 路由前缀 /system/tasks,非 AIP 的 /tasks30s;响应 data 含 status(running/success/failed/cancelled)、progressmessageerrorresult
GET/pipelines/:id/runsquery: limit/offset(0=不分页,缺省全量,id 降序)30s
POST/pipelines/:id/enable30s
POST/pipelines/:id/disable30s

5.3 响应结构

统一成功包装 {code: 0, data: ...}okData);失败 {code, error}errorResponse)。

创建成功:

{ "code": 0, "data": { "id": 3 } }

管道列表(GET /pipelines):

{
  "code": 0,
  "data": [
    {
      "id": 3,
      "name": "orders_sync",
      "description": "",
      "source_connector_id": 1,
      "target_object_type_id": 2,
      "sql_steps": [ { "name": "extract", "sql": "SELECT * FROM orders", "transform_type": "select" } ],
      "write_mode": "append",
      "trigger_type": "cron",
      "cron_expr": "0 2 * * *",
      "enabled": true
    }
  ]
}

运行触发(POST /pipelines/:id/run):

{
  "code": 0,
  "data": {
    "id": 17,
    "pipeline_id": 3,
    "status": "running",
    "started_at": "2026-08-30T02:00:00Z",
    "rows_in": 0,
    "rows_out": 0,
    "triggered_by": "web"
  }
}

运行历史(GET /pipelines/:id/runs):

{
  "code": 0,
  "data": [
    { "id": 17, "pipeline_id": 3, "status": "success", "started_at": "...", "finished_at": "...", "rows_in": 1000, "rows_out": 1000, "error": "", "triggered_by": "web" }
  ]
}

5.4 关联模块表

后端包/文件职责
products/foundry/server/pipeline_handlers.go管道 CRUD/运行/启停/运行历史 handler(同文件还有血缘与质量 handler)
products/foundry/pipeline/models.goPipelineDefinition / PipelineRun / DataQualityRule / DataQualityIssue(表 pipeline_definitions / pipeline_runs / data_quality_rules / data_quality_issues)
products/foundry/pipeline/service.goCRUD 校验、RunPipeline 七步执行、write_mode 同步、血缘打点、调度器同步
products/foundry/pipeline/quality.go质量规则检查与问题记录(null/format/unique/referential)
products/foundry/lineage血缘打点(RecordSourceToPipeline / RecordPipelineToObject)
platform/scheduler统一调度器(V5 Stage 0 注入,cron_expr 激活)
platform/connector真实数据源连接器(DefaultConnectorProvider 白名单工厂)
products/foundry/ontology目标对象类型读取(ObjectTypeReader)

5.5 关键机制

6. 核心流程详解

6.1 创建一条管道(主流程)

  1. 点「新建管道」展开表单;
  2. 填名称(api_name)、源连接器 ID(数据源 id,必填)、目标对象类型 ID(可选);
  3. 选写入模式(默认 append)与触发类型(默认 manual;选 cron 则补 cron 表达式);
  4. 填描述、保持「启用」勾选;
  5. 「SQL 步骤(至少 1 个)」区点「+ 步骤」,填写步骤名、选变换类型(select/insert/update/create_table_as)、写 SQL;
  6. 点「保存」:前端校验——过滤出「名称与 SQL 均非空」的步骤,无有效步骤则 alert「至少需要 1 个有效 SQL 步骤(名称与 SQL 均需填写)」;cron 类型无表达式则 alert「cron 触发类型需要填写 cron 表达式」;
  7. 创建走 POST /pipelines(附 name),成功 alert 管道创建成功 id=<id> 并刷新列表;后端同步注册 cron 调度。

6.2 编辑管道

6.3 运行与终态判定(后端任务化 + 前端退避轮询)

  1. 点「运行」→ confirm → POST /pipelines/:id/run(body {triggered_by:'web'});
  2. 后端任务化双响应:
    • 快路径 200:15s 内执行完成,data 为 run 记录(含 task_id),alert 管道运行已启动 run_id=<id> status=<status>
    • 超时受理 202:执行超过 15s 转入后台,data={task_id,status:"running",type:"pipeline_run"},alert 管道已受理,正在执行(task_id=…)并启动轮询。
  3. 轮询终态(退避 + 封顶 + 超时后继续轮询):轮询 GET /system/tasks/:id,延迟 1.5s 起、每次 ×1.5、封顶 10s,总上限 5 分钟,直至:
    • success:alert 管道运行成功 run_id=…(run 记录在 task.result),已展开该管道历史则自动刷新;
    • failed/cancelled:alert-error 展示 task.error/task.message,并自动刷新运行历史;
    • 单次轮询请求失败不中断任务,按退避节奏继续重试;
    • 到达 5 分钟上限仍未终态:alert 提示「轮询已达 5 分钟上限;可点击「继续轮询」获取终态」,并在该管道行露出「继续轮询」按钮,由用户一键续轮询直至终态。
  4. 轮询期该管道「运行」按钮显示「运行中」并禁用;组件卸载自动停止轮询(onBeforeUnmount(stopRunPolling) 清定时器并置取消标记)。该模式与「实体消解」(FusionPage.vue)的运行轮询一致。
  5. 成功路径:rows_in/rows_out 反映实际读入/写出行数,且自动打血缘(可在「数据血缘」页按管道名查上游/下游)。
提示:POST /run 已任务化(15s 快路径 200 / 超时 202 task_id),页面会自动退避轮询至终态;若轮询达 5 分钟上限,点该管道行「继续轮询」续轮询,不要重复点「运行」制造多份 run。

6.4 启用/停用与删除

6.5 cron 定时调度流

7. 权限与安全

8. 常见问题与排错

8.1 点「运行」后 alert「管道保存失败」或「运行失败」且无明确错误

8.2 运行历史一直停在「running」或状态不更新

8.3 保存管道报「pipeline "xx" already exists」或重复键错误

8.4 cron 管道不按时触发

8.5 保存/运行管道报「write_mode=changelog 暂未实现」

9. 已知缺陷与边界

缺陷/边界说明影响
changelog 写模式 P2(设计如此,非缺陷按主键 upsert + 变更历史未实现(W2-01 起改为显式拒绝以免静默按普通写入误导用户)创建/更新即 400 拒绝(pipeline/service.go:415-416),前端下拉禁用,不进入可运行态(存量数据运行时仍有兜底标 failed)
跨数据源不支持(F3)管道步骤限定单数据源(source_connector_id 一个)无法跨源 JOIN 组装
运行异步受理POST /run 15s 内完成返回 200(run 信息),超时 202 task_id 转后台执行本页已适配:自动退避轮询 GET /system/tasks/:id 直至 success/failed/cancelled 并提示/刷新历史
轮询等待上限 + 续轮询(已于 2026-09-13 升级)单次运行自动轮询 1.5s 起 ×1.5 封顶 10s、总上限 5 分钟达上限后自动轮询暂停并露出「继续轮询」按钮,用户一键续轮询直至终态(不再停留在「请自行去运行历史确认」)
update 步骤 rows_in 不可靠DDL 的 RowsAffected 在 SQLite 恒为 1,不计入 rows_in行数统计近似
血缘打点非阻断血缘失败/未注入不影响 run 结果血缘可能缺失
质量检查失败跳过CheckQuality 出错时 continue(不阻断 run)质量检查偶发失败无告警
定时调度依赖注入scheduler 为 nil 时 cron 不激活需 V5 Stage 0 接线注入
snapshot 无备份全量替换先 DELETE 再 INSERT,无事务回滚中断可能留空表