# 数据管道：从源系统到对象

对象背后是一条数据管道：**源头只读一次，只写变化，有人要的时候才算，健康状况看得见**。这一页按数据流动的顺序讲清每一段由谁定义、在哪里看、有什么规则，再算一笔「定时刷新 vs 按需」的计算账，最后给出把一个定时看板迁进 Semantic 的步骤。对象本身怎么定义见 [Object](objects.md)，全部写法见 [Semantic](semantic.md) 与 [闭环](loop.md)。示例取自虚构的 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 页改增量、挑列。

```bash
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

```bash
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 行；同一个导入同时只有一个在途的构建。
- **数据边界**：数据集的字节在这家公司自己的存储桶里，平台库里还存对象属性值（同步来的值和动作改过的值）；配置、查询与数据集只给本组织开发者看。

```bash
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 秒内算完；一个数据集只有一个产出者。

```json
{ "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"] } }
```

```bash
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`** 里——改数据源就是改本体。

```json
"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 的数据流——不要让两个数据源映射同一个属性；有多个数据源的类型不能配对象 / 属性安全策略。

```json
"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" } }
]
```

```bash
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](objects.md)）。
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`。

完整的自动化写法（条件、效果、重试、失败效果）与停用的旧效果怎么替换，见 [闭环](loop.md) 第 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 分钟断开，切回续上；连接时长计入应用的计算分钟 |
