查看 Markdown

数据管道:从源系统到对象

对象背后是一条数据管道:源头只读一次,只写变化,有人要的时候才算,健康状况看得见。这一页按数据流动的顺序讲清每一段由谁定义、在哪里看、有什么规则,再算一笔「定时刷新 vs 按需」的计算账,最后给出把一个定时看板迁进 Semantic 的步骤。对象本身怎么定义见 Object,全部写法见 Semantic 与 闭环。示例取自虚构的 Demo Company。

数据源(ERP / MES / OA,只读)
   │  ① 接入:数据流(发布端只发差)或 Data Connection(平台用只读 SQL 读,客户内网经 agent proxy)
   ▼
Dataset(每次导入一个事务;字节在这家公司自己的存储桶里)
   │  ② 加工:Transform 从数据集算出新的数据集(只重算受影响的部分)
   ▼
Object Type 的数据源(propertyMapping)── 只把变了的行写成对象,源里没了的标「源头已消失」
   │  ③ 读:应用、智能体、SQL、订阅;数据比新鲜度要求旧才按需同步
   ▼
Automation(数据一变 / 到点)──▶ Action · Function · 邮件

1. 六段,各在哪里看

段 做什么 谁定义 在哪里看(/semantic/<组织>)
数据源 公司原有的系统,永远只读 客户的系统 —
接入 数据流(发布端推差)或 TableImport(平台读) 开发者 Data Connection:Sources、Agents、batch sync(Runs)
数据集 一次导入或加工的产出,按事务提交 平台 Data Connection:dataset 页(事务历史、Health)
加工 从数据集算出新的数据集 开发者(定义是 JSON) aidc semantic transforms show
对象 数据集 / 数据流映射到对象类型的属性 本体(对象类型定义的 datasources) Ontology Manager 的数据源页、Object Explorer
使用 应用、智能体、SQL、订阅、自动化 应用与自动化的作者 Applications、Automate 的运行记录

2. 接入:两条路

数据流(stream) Data Connection(TableImport)
谁去读数据源 装在数据源旁边的发布端(aidc semantic streams pipe + 适配器,常驻服务,不是定时器) 平台用只读账号执行一条 SELECT
发什么 只发与上次相比的差 每次执行写一个数据集事务,平台再与上一份比差
凭证 发布 Key,只能往这一条流发布 数据库口令写一次,平台加密保存、不再回显
客户内网 发布端本来就在内网 经装在客户网络里的 agent proxy(只出站、只转发字节)
适合 能在数据源旁边装程序、要秒级推送的数据(MES 报工) 不想在数据源旁边装程序的 ERP / MES 表

两条路都只读:平台永远不写数据源,不另开备份库。Demo Company:MES 报工走数据流 mes-reports,ERP 的订单、工单、物料、库存流水、采购走 Data Connection。

先探索源,再建同步

source 页的 Explore 列出源里的表与视图(可按表名搜索)、一张表的列、主键、外键与被引用的表;直连的 source 能取几行样本(样本不落库),经 agent proxy 的 source 只看结构。勾选要的表一次建好 batch sync(每张一条 SELECT *),再到 sync 页改增量、挑列。

aidc semantic connectivity explore <connectionRid> --search order
aidc semantic connectivity explore <connectionRid> --table dbo.SalesOrderLine
aidc semantic connectivity explore <connectionRid> --table dbo.SalesOrderLine --preview --rows 5

只选用得到的表和列,查询带时间窗口(例:近 90 天)——读得少,后面每一段都省。

客户内网:agent proxy

aidc semantic connectivity agent register --ontology cell-demo --api-name site-box --display-name "工厂机房"   # 凭证只给一次 + 三步安装
aidc semantic connectivity egress create --ontology cell-demo --file erp-via-agent.json                  # ERP 的主机:端口,经 site-box
aidc semantic connectivity connection --ontology cell-demo --file erp-connection.json                    # JDBC 地址 + 只读账号 + 出口策略
  • 查询在平台的 worker 里执行;agent 只为这次构建开一条隧道、转发字节,不保存数据。
  • agent 与平台之间的控制通道常驻(平台随时能请它开隧道);隧道只在构建或探索源的时候开。
  • 换口令从标准输入读,不进命令行:printf '%s' "$PASSWORD" | aidc semantic connectivity secret <connectionRid> Password。

3. 数据集与事务

每个 TableImport 执行一次 = 数据集的一个事务;平台与上一份比差,只把变了的行写成对象,没变的一行都不写。

导入模式 对象怎么变 适合
SNAPSHOT(可带窗口查询) 对象就是这次读到的行;窗口外的、源里没了的标「源头已消失」 订单、工单这类「只要最近一段」的数据
APPEND + 增量列(? 绑上次水位) 只查上次之后改过的行,只增改 有可靠修改时间的流水表(库存流水 tUpdateTime)
  • 每周对账:APPEND 看不到源头删掉的行,平台大约每周一次对增量导入做全量快照,源头删掉的行在这时标「源头已消失」(距上次全量满 7 天、最近 7 天有成功运行的导入才做)。
  • 上限:每个 TableImport ≤ 150,000 行;同一个导入同时只有一个在途的构建。
  • 数据边界:数据集的字节在这家公司自己的存储桶里,平台库里还存对象属性值(同步来的值和动作改过的值);配置、查询与数据集只给本组织开发者看。
aidc semantic connectivity import <connectionRid> --file sales-order-lines.json    # SNAPSHOT + 窗口;或 APPEND + 增量列
aidc semantic connectivity execute <connectionRid> <tableImportRid> --wait         # 第一次全量,之后只写变化

4. 加工:一次读源,多次加工

指标不要写在导入的 SQL 里、每次回源系统整份算:原始数据进数据集只读一次,指标由 Transform 从数据集算成新的数据集(一条 DuckDB 的 SELECT,输入按别名当表用),再作对象类型的数据源。输入一提交新事务就自动构建;同一个 Transform 同时只跑一个。

写法 每次算什么 适合
非增量(不写 incremental) 整份重算 小表;要看前后行(排名、前一天的值)
append 只算新加的行,追加到产出 逐行的清洗、过滤、补字段
mergeAndReplace 新行影响到的键在输入里全部重算,再按键合并 按键汇总(按物料、按日期 × 车间)
snapshotInputs 参考表整份读,它的变化不触发 物料清单、价格表

整份重算只在三种时候发生:第一次构建、semanticVersion 变了(改了口径要回溯历史就加一)、非 snapshot 输入被整份替换。requireIncremental: true 时,除了第一次构建和 semanticVersion 变了,不能增量就失败,防止大表意外整份算。边界:一次构建读的文件 ≤ 384 MB、产出 ≤ 500 万行、300 秒内算完;一个数据集只有一个产出者。

{ "apiName": "stock-balance", "displayName": "库存余额",
  "inputs": {
    "moves": { "dataset": "stock-moves", "key": ["material_code"] },
    "materials": { "dataset": "materials" } },
  "sql": "SELECT m.material_code, sum(m.qty) AS on_hand, max(m.moved_at) AS last_move FROM moves m SEMI JOIN materials t ON t.material_code = m.material_code GROUP BY ALL",
  "incremental": { "semanticVersion": 1, "snapshotInputs": ["materials"], "output": "mergeAndReplace", "key": ["material_code"] } }
aidc semantic transforms put --ontology cell-demo --file stock-balance.json --dry-run   # 先预演:输入、环、名字
aidc semantic transforms put --ontology cell-demo --file stock-balance.json
aidc semantic transforms build <transformRid> --wait
aidc semantic transforms show <transformRid>                                              # 增量还是整份、受影响的键、读写多少行

5. 接到对象类型

数据集或数据流接到哪个对象类型、哪一列对哪个属性,写在对象类型定义的 datasources 里——改数据源就是改本体。

"datasources": [
  { "type": "dataset", "dataset": "stock-balance", "propertyMapping": { "materialCode": "material_code", "onHand": "on_hand" } }
]
  • 一个对象类型恰好一个主键属性,主键必须映射;属性 API name 用 camelCase,列名照源系统写。
  • 只在语义层维护的属性(writeback)和派生属性不映射;同一条流在一个类型里只出现一次。
  • 数据流的 mode:mirror = 流里没了的行标「源头已消失」;upsert = 只增改。
  • 重新定义时没写 datasources = 沿用现有的(返回警告);要去掉全部,写 "datasources": []。
  • 定义落了之后先全量同步一次,之后增量。

一个对象类型可以有多个数据源,每个数据源都要映射主键。Demo Company 的工单:计划来自 ERP 的数据集,完成数量来自 MES 的数据流——不要让两个数据源映射同一个属性;有多个数据源的类型不能配对象 / 属性安全策略。

"datasources": [
  { "type": "dataset", "dataset": "work-orders", "propertyMapping": { "workOrderNo": "wo_no", "lineCode": "line_code", "plannedStart": "planned_start", "qtyPlanned": "qty_planned", "status": "status" } },
  { "type": "stream", "stream": "mes-reports", "mode": "upsert", "propertyMapping": { "workOrderNo": "wo_no", "qtyDone": "qty_done" } }
]
aidc semantic datasource list      # 数据流的绑定与同步进度
aidc semantic datasource sync      # 平台数据源立即同步(= /semantic 的「立即同步」)

6. 新鲜度与按需

读由数据源同步的类型时,数据比 freshness.maxAgeSeconds(缺省 300 秒)旧才按需同步一次:先返回当前值,后台同步;?fresh=true(Python 客户端 fresh=True)最多等 20 秒,到时还没同步完就返回现有的数据。没人读、也没有自动化引用 = 不同步 = 不花钱。

什么在消耗计算 什么时候发生
数据源的同步 有人读且数据比 maxAge 旧;或被自动化的对象集条件引用(按评估频率,≥ 5 分钟,缺省 10 分钟,可限工作时段)
加工 输入提交了新事务
应用 有人打开;只算当前页
订阅 页面可见时连着;隐藏超过 1 分钟断开,切回续上
自动化 数据一变(对象集条件)或到点(时间条件)

Data Connection 这条路上,唯一常驻的是 agent proxy 的控制通道;只用直连的组织连这个也没有。数据流的发布端装在数据源旁边,常驻运行、定时发心跳。

7. 数据健康

数据集上可以挂检查(Data Connection 的 dataset 页右侧「Health」),失败、升级、恢复时给建检查的人发邮件(同一状态 6 小时内不重发;报告与邮件只有状态和计数,不写数据)。

检查 怎么判 什么时候评估
buildStatus 最近一次构建成败(每周对账的那次不算);可设连续失败几次升级成 CRITICAL 有构建结束时
timeSinceLastUpdated 距最近一个有新行的事务多久 定期
schemaComparison 数据集现在的列与期望的列比 有新事务时
primaryKey 这几列非空、不重复 有新事务时

按需运行意味着没人看的时段数据不会更新:timeSinceLastUpdated 的上限按最长的正常空闲设(例如每天都有人读的数据设 26 小时;周末没人读的,要盖住整个周末);要发现「有人要数据却拿不到」,看 buildStatus(构建失败、排队超时都会报)。

8. 计算账:定时刷新 vs 按需

旧做法里,计算量由时钟决定;按需之后,计算量由有人读、数据变化决定。

每天运行次数 = 定时任务数 × (1440 ÷ 间隔分钟)
每天计算量  ≈ 每天运行次数 × 每次耗时(× 每次读的数据量)
Demo Company(示例) 定时刷新 按需
24 个看板脚本,每 10 分钟一次 24 × 144 = 3,456 次 / 天,每次查 ERP 加生成页面约 20 秒 → 每天约 19 小时计算 看板打开时现算当前页;数据过期才同步,同一类数据一小时最多 12 次,只在有人读的时段
临时隧道的保活任务,每分钟一次 1,440 次 / 天 应用发布在平台上,没有隧道、没有保活
缺料提醒,每 5 分钟查一遍 ERP 比快照 288 次 / 天,大多数次什么都没变 自动化的对象集条件:对象有变化(同步进来、动作写入)才评估;按时拉数据源可限工作时段(workHours)
数据是不是新的 脚本断了也没人知道 类型带同步状态,应用有数据新鲜度组件;数据健康检查失败会告警

这一类账在现场反复出现:服务器上的定时任务越积越多,很多看板没人看也每隔几分钟跑一次。按「数据一变、到点、有人看」三种触发重写之后,绝大多数可以停掉。另外两条规则:大表(例如 2 万行以上)的首次全量放到夜里、写入时盯错误率;平时只写变化。

9. 迁移:把一个定时看板搬进 Semantic

  1. 清点:每个定时任务的频率、读什么、产出什么、谁在看、每天真正被打开几次。
  2. 分类:数据一变才要做的(→ 自动化的对象集条件);时间本身是条件的(日报、月底对账 → 时间条件,每天一次是常态);有人看才需要的(→ 应用按需);没人用的(→ 直接停)。
  3. 建模:按业务对象建对象类型,不镜像原表(Object)。
  4. 接数据:Explore → TableImport(SNAPSHOT + 窗口或 APPEND + 增量列)或数据流 → 需要的加工 → datasources → 新鲜度。
  5. 做应用或自动化:看板做成 Applications(只存定义,打开时现算);提醒做成自动化,写上 replaces(替掉了哪个任务、它原来每天跑几次)。
  6. 预演:aidc semantic automations upsert … --dry-run,看编译出来的对象集与校验。
  7. 并行一轮:新旧同时跑,对上结果。
  8. 暂停旧任务:暂停不删,保留两周,留好恢复的办法;关掉隧道与保活。
  9. 记账:自动化的详情页(/semantic/<组织>/automate/automations/<apiName>)「替代的 cron」一栏,列出它替掉了哪些定时任务、原来每天跑多少次。
  10. 挂健康检查:在关键数据集上加 buildStatus 与 primaryKey。

完整的自动化写法(条件、效果、重试、失败效果)与停用的旧效果怎么替换,见 闭环 第 9 节。

10. 上限一览

项 上限 / 缺省
每个 TableImport ≤ 150,000 行;同时一个在途构建
新鲜度 freshness.maxAgeSeconds 缺省 300 秒;fresh=true 最多等 20 秒
自动化评估频率 evaluation.everyMinutes ≥ 5,缺省 10
Transform 一次构建 读 ≤ 384 MB、产出 ≤ 500 万行、≤ 300 秒
数据流 单个事件 ≤ 1 MB、最新状态 ≤ 4 MB、每批 ≤ 100 条、每条流每分钟 ≤ 600 次发布
订阅 页面隐藏超过 1 分钟断开,切回续上;连接时长计入应用的计算分钟

本页由 developer/docs/pipeline.md 生成 · Markdown 原文 · llms.txt