P2 Foundry
同步引擎:全量与增量
把外部数据源按计划持续抽进平台内数据集:全量(full)在单事务内清空重写,增量(incremental)用水位线断点续传、幂等键冲突跳过;运行经任务系统异步化可看进度,成功后血缘打点与质量画像自动刷新。看完这 3 个故事,你就能把“每天自动更新数据”这件事声明式地管起来。
数据工程师
运维人员
全量重写
增量水位
定时调度
任务进度
共 3 个故事
能 / 不能速览
✅ 这个主题能做
- 全量(full):单事务内 DELETE 旧行 + INSERT 新行,原子重写目标数据集
- 增量(incremental):watermark 水位线推进 + 断点续传,field_map 幂等键冲突跳过
- 运行任务化(type=sync_run),前端轮询任务进度条,运行历史保留成功 / 失败记录
- schedule cron 经统一调度器注册(job 名 sync:<id>),定时触发照常落运行记录
- 成功链:血缘打点(源 → 数据集)+ 质量画像自动刷新,best-effort 失败不阻断
⛔ 这个主题做不了
- object 来源目标创建可用,但运行时明确报错(对象行读取器未注入)
- 增量只追加:源端删除的行不会反映到数据集,需 full 周期重写兜底
- 运行失败无自动重试;同一目标无防重入(手动 + cron 同时触发可能并发执行)
- sync_run 的 runner 未注册任务系统,重启后 queued 任务按默认语义标 failed(可手动重跑)
- full 无数据量上限,超大源会长时间占事务;来源查询建议带 LIMIT
适用角色
本主题面向两个角色:
- 数据工程师:建同步目标、选模式、配水位 / 幂等键、排错重跑——是核心使用者。
- 运维人员:通过运行历史与调度任务核对数据健康度,确认定时任务是否在跑。
下游分析师消费目标数据集;平台审计里能看到每次运行的触发人 / 触发来源(manual / cron)与结果。
能力速览(能做什么)
全量重写
full 模式在单事务内 DELETE 旧行 + 逐行参数化 INSERT,失败整体回滚、旧数据不动;row_count 以行表 COUNT(*) 为准。
增量水位
incremental 把查询包装为“水位过滤 + 参数绑定”,首次全量建立水位,之后只拉新增行;水位值存目标里可断点续传。
幂等键防重
field_map 目标列作为业务幂等键,先读行表键集合、冲突行(含批内重复)跳过,rows_synced 不含跳过行。
任务化与进度
运行经 task.Manager 任务化(type=sync_run)返回 task_id,前端 1.5s 轮询进度条;任务系统不可用自动降级同步执行。
定时调度
enabled + schedule 目标经统一调度器注册(sync:<id>),cron 触发 triggered_by=cron 并照常落运行记录,启动时重放恢复。
成功链
同步成功后自动血缘打点(源 → 数据集)+ 触发该数据集启用的质量画像,“数据进来即体检”,均失败仅记日志不阻断。
调整指南(怎么调整)
- 选模式:数据量小 / 源表重写频繁 / 需要源端删除同步 → full(事务内重写,简单可靠);数据量大且只追加 → incremental(水位推进 + 幂等键防重)。两者可混用:日常 incremental,周期 full 兜底。
- 配水位:watermark_col 填源查询结果中单调递增的列(时间戳 / 自增 id);incremental 必填且须为目标结果列之一,full 可选(记录了为将来切 incremental 预留连续起点)。
- 配幂等键:field_map 的“目标列”就是增量幂等键(业务主键语义);不配 field_map 时列按归一名自动匹配,防重仅靠水位过滤。
- 改调度:enabled=false 即停用(停用后 /run 返回 400),创建时非法 cron 会回滚删除目标行;启动时会重放存量调度,重启后定时任务自动恢复。
- 改来源:来源查询建议带 LIMIT 或分区;字段名拼错导致失败后,改定义手动重跑即可。
做得好的场景
同步引擎把“外部数据持续入平台”变成声明式、可调度、可追溯的抽取流程,特别适合以下场景:
- 例行数据刷新:每天 2 点 / 每 2 小时 cron 自动拉数,代替手工导表,漏跑有运行历史可查。
- 大表增量抽取:水位过滤只拉新增行,避免全表反复重导;断点就在 watermark_value,重启 / 失败后重跑即续传。
- 防重入与防重:cron 侧 SkipIfStillRunning 双保险,幂等键业务层去重,重复数据进不了数据集。
- 数据进来即体检:同步成功自动触发质量画像,坏数据在源头就被发现,而不是等下游报表出错。
限制与不足
以下是明确的边界,使用前先知道:
- object 源未接线:source_type=object 的目标创建可用,运行时明确报错 503(读取器未注入)。
- 增量不删:watermark 增量只追加,源端删除的行不会反映,需 full 周期重写或后续 CDC。
- 重启恢复语义:sync_run runner 未注册任务系统,重启后 queued 任务标 failed,但 fs_sync_runs 保留、可手动重跑。
- 失败无重试:运行失败即 failed,无自动重试;同一目标手动 + cron 同时触发可能并发执行。
- 幂等键非强约束:行表无唯一索引,“先查后插”在单任务内原子,跨任务并发同键仍可能重复。
场景故事
故事 1
订单快照每 2 小时全量重写:任务化 + 进度轮询 + 成功链
场景:全量同步
角色:数据工程师
耗时:约 8 分钟
- 背景
- 老赵想让“订单快照”每 2 小时自动全量刷新。他先建好目标数据集 orders_snapshot,再在同步工作台建目标 orders_hourly:source_type=platform_ds、ds_id=1、来源查询取订单核心列、mode=full、schedule=0 */2 * * *,创建后立即运行一次验证。
- 传统做法对比
- 以前全量刷新要么手工导表(漏一次报表就错一天),要么写 crontab + 脚本(部署半天、跑没跑全凭日志);现在一条目标 = 声明式配置,运行记录、进度、成功链全自动,前端还能看进度条。
- 角色
- 数据工程师(建目标、看运行);运维人员通过运行历史与调度任务核对健康度。
- 操作步骤
-
- “数据同步”页新建目标:名称 orders_hourly、目标数据集 orders_snapshot、来源 platform_ds、ds_id=1
- 来源 SQL 填 SELECT order_id, amount, status, updated_at FROM orders
- 模式选 full,调度填 0 */2 * * *,启用勾选
- 点“立即运行”,前端出现任务进度条(1.5s 轮询)
- 看运行历史确认 status=success、rows_synced=5
- 系统响应
- 运行与进度返回示例:
POST /api/v1/sync/targets/t1-.../run
→ {"code":0,"data":{"run":{"id":"r1-...","status":"running","mode":"full","task_id":"task-abc"},"task_id":"task-abc"}}
GET /api/v1/system/tasks/task-abc
→ {"code":0,"data":{"id":"task-abc","type":"sync_run","status":"success","progress":1,"message":"血缘与质量画像","result":{"run_id":"r1-...","status":"success","rows_synced":5}}}
- 结果洞察
- full 模式在单事务内 DELETE 旧行 + INSERT 新行,第二次运行后行表只有最新数据,中途失败整体回滚、旧数据不动;row_count 以行表 COUNT(*) 为准,size_bytes 取本次写入。成功后血缘边 datasource:1 → dataset 与质量画像自动触发(best-effort,失败不阻断同步成功)。目标 last_run_at 成功 / 失败都更新,运行痕迹诚实记录。
- 调整建议
- 来源查询建议带 LIMIT 或分区,full 是单事务逐行写,超大源会长时间占事务;源端有删除需求就用 full 周期兜底(incremental 只追加);想更细看进度就关注任务阶段消息(加载数据集 → 拉取源数据 → 写入 → 更新元数据 → 血缘与画像)。
- 动手试一试
- 登录:admin / admin1。页面路径:数据同步 → 新建。输入内容:名称 orders_hourly、数据集 orders_snapshot、ds_id=1、SQL SELECT order_id, amount, status, updated_at FROM orders、mode=full。预期结果:立即运行返回 task_id,轮询后 success、rows_synced=5。
- 限制提示
- full 无数据量上限,超大源注意长事务;同一目标无防重入(手动 + cron 同时触发可能并发执行,cron 侧有 SkipIfStillRunning);object 来源仅支持 full 且需接线注入读取器,当前未注入运行时报 503。
故事 2
订单增量同步:水位线断点续传 + 幂等键冲突跳过
场景:增量同步
角色:数据工程师
耗时:约 8 分钟
- 背景
- 订单量上来后,老赵把同步目标改成增量:watermark_col=updated_at、mode=incremental、field_map={“order_id”:“order_id”}。首次运行全量拉取建立水位,之后每次只拉 updated_at 大于上次水位的行追加,业务键冲突自动跳过。今天源表新增了 3 条订单,他要验证二次运行。
- 传统做法对比
- 以前增量要么全表重导(越导越慢),要么自己写“取大于某时间戳”的 SQL 定时跑(水位存哪、断了怎么续全靠脑记);现在水位值存在目标里(watermark_value),断点续传自动完成,field_map 幂等键再兜一层防重。
- 角色
- 数据工程师(设计增量键与水位列);下游分析师关注 rows_synced 与数据时效。
- 操作步骤
-
- 新建目标:名称 orders_incremental、数据集 orders_snapshot、来源 platform_ds
- 来源 SQL 填 SELECT order_id, amount, status, updated_at FROM orders,watermark_col=updated_at
- 模式选 incremental,field_map 填 {“order_id”:“order_id”}
- 首次运行(无水位 → 全量拉取 → 建立水位值)
- 往源表插入 3 条 updated_at 更新的订单后再次运行,观察 rows_synced=3、水位推进
- 系统响应
- 二次运行与水位返回示例:
POST /api/v1/sync/targets/t1-.../run
→ {"code":0,"data":{"run":{"id":"r2-...","status":"running","mode":"incremental","task_id":"task-def"},"task_id":"task-def"}}
GET /api/v1/sync/targets/t1-...
→ {"code":0,"data":{"target":{"mode":"incremental","watermark_value":"2026-08-30T10:00:00+08:00","last_run_at":"2026-08-30T10:00:01+08:00"}}}
GET /api/v1/sync/targets/t1-.../runs → 首次 rows_synced=5,二次 rows_synced=3
- 结果洞察
- 增量二次运行把查询包装成 SELECT * FROM (SELECT ...) WHERE updated_at > ?(水位值参数绑定,不内联拼 SQL),只追加新行;field_map 幂等键(order_id)先读行表已有键集合,冲突行(含批内重复)跳过,rows_synced 不含跳过行。水位按数值取最大,避免 “9” > “10” 这类文本比较失真。断点就在 watermark_value 里,失败后手动重跑即续传。
- 调整建议
- watermark_col 选单调递增列(时间戳 / 自增 id);幂等键取 field_map 的目标列(业务主键语义),不要默认用水位列,否则同水位不同业务键的合法新行会被误丢;增量不删源端删除的行,需删除同步配 full 周期重写;大行表上幂等键会先全表读键集合,建议定期 full 压缩。
- 动手试一试
- 登录:admin / admin1。页面路径:数据同步 → 新建。输入内容:名称 orders_incremental、mode=incremental、watermark_col=updated_at、field_map={“order_id”:“order_id”}。预期结果:首次运行 rows_synced=5 并建立水位;向源表插 3 条新订单后再跑,rows_synced=3。
- 限制提示
- incremental 只追加,不反映源端删除;水位存文本、绑定 SQL 时按内容转数值,日期型水位须保证可比较(RFC3339 可用);行表无唯一索引,跨任务并发同键写入仍可能重复(单任务内原子);源查询结果必须含 watermark 列,否则 422。
故事 3
凌晨 2 点跑挂了:从运行历史定位,改定义手动重跑,调度痕迹可查
场景:调度与排错
角色:数据工程师
耗时:约 8 分钟
- 背景
- orders_hourly 每 2 小时 cron 触发一次。某天早上老赵发现最新运行 failed,error 指向来源查询字段拼写错误;他修正目标定义后手动重跑成功,并到调度器核对 job 名 sync:<targetID> 的注册与触发痕迹。
- 传统做法对比
- 以前定时任务挂在服务器 crontab 上,跑挂没人知道,日志分散难排查;现在调度条目注册在统一调度器,运行痕迹落在 scheduler_runs,运行历史里失败记录与 error 一起保留,改完即重跑,全程有据可查。
- 角色
- 数据工程师(排查与重跑);运维人员通过调度任务列表看触发痕迹。
- 操作步骤
-
- 数据同步列表看到最近运行 failed,点“运行历史”看 error
- 打开调度器任务页,确认 job 名 sync:<targetID> 与 cron 已注册
- 编辑目标修正来源 SQL(如字段名拼错)
- 重新“立即运行”,轮询任务到 success
- 回运行历史确认 triggered_by=manual、rows_synced 正常
- 系统响应
- 运行历史与调度任务返回示例:
GET /api/v1/sync/targets/t1-.../runs
→ {"code":0,"data":{"runs":[
{"id":"r3-...","status":"failed","rows_synced":0,"error":"SQL_EXECUTION_ERROR: ...no such column: updated_at2","triggered_by":"cron"},
{"id":"r4-...","status":"success","rows_synced":5,"triggered_by":"manual"}]}}
GET /api/v1/scheduler/jobs → 可见 job name "sync:<targetID>" 与其 cron spec
- 结果洞察
- 调度与运行分离:cron 触发走内部同步执行(triggered_by=cron,不经任务系统,保持调度器防重入语义),运行痕迹照常落 fs_sync_runs 并可经调度器查看。失败不可怕——error 精确可查、运行记录保留、修正后手动重跑即恢复。要注意重启语义:sync_run 的 runner 未注册任务系统,重启后 queued 任务按默认语义标 failed,但 fs_sync_runs 保留,可手动重跑补齐。
- 调整建议
- 排查顺序:运行历史 error → 目标定义 → 源表结构;改完定义建议先手动跑一次再等 cron;server 启动时会重放存量 enabled+schedule 目标注册调度,重启后定时任务自动恢复;停用目标(enabled=false)后 /run 返回 400,手动运行前先启用。
- 动手试一试
- 登录:admin / admin1。页面路径:数据同步 → 运行历史。输入内容:把目标来源 SQL 某字段名故意写错后立即运行。预期结果:run 返回 failed、error 含 no such column;改回正确列名重跑后 success,历史保留两条记录。
- 限制提示
- 调度条目在创建 / 更新时同步注册,删除目标会反注册(失败仅告警);object 来源目标创建可用但执行必失败(读取器未注入);运行失败无自动重试;质量画像的调度仅启动时批量注册,改 schedule 需重启才生效。
常见问题
full 和 incremental 怎么选?
数据量小 / 源表重写频繁 / 需要源端删除同步 → full(事务内重写,简单可靠);数据量大且只追加 / 按时间增量 → incremental(水位推进 + 幂等键防重)。两者可混用:日常 incremental,周期 full 兜底。
watermark_col 到底填什么?
源查询结果中单调递增的列(时间戳 / 自增 id)。incremental 模式必填且须为目标结果列之一;full 模式可选(配置了且结果含该列则记录水位,为将来切 incremental 预留连续起点)。
为什么增量用水位过滤还不够,还要幂等键?
水位 > 过滤保证“上次水位之后”的行只取一次,但同水位内的合法重复(如业务键相同的新行)会被误丢或误重;field_map 的幂等键在业务层去重,双保险。
同步失败会影响数据集已有数据吗?
full 模式在单事务内 DELETE + INSERT,失败整体回滚、旧数据不变;incremental 追加也在事务内,事务失败即整体回滚。均以 fs_sync_runs.error 为准排查。
重启后定时任务还能跑吗?
启动时会重放存量 enabled + schedule 目标的调度注册,定时任务自动恢复;但重启时 queued 的 sync_run 任务因 runner 未注册会标 failed,运行记录保留、可手动重跑。
为什么 object 来源同步会报错?
object 来源需要平台注入“对象行读取器”(ObjectRowReader),当前接线处未注入,运行时明确报 503。object 目标可创建但执行必失败,属诚实降级,读取器接线在后续版本补齐。
主题小结
一句话:同步引擎把“外部数据持续入平台”变成声明式、可调度、可追溯的抽取流程——full 事务内重写、incremental 水位断点续传 + 幂等键防重、任务化可看进度、成功链自动接血缘与画像。记住几个边界:object 源未接线、增量不删源端删除、失败无自动重试、重启后 queued 任务标 failed 可手动重跑。