数据源只读地进来

第 1 课 · 共 6 课 约 9 分钟

MES、ERP 等数据源经发布端、数据流和同步绑定进入对象类型;接入只读源系统,Action 可配置 writeback webhook;主键要唯一而且稳定。

本课目标

读完这一课,你将能够

  • 说出数据从示例工厂的 MES 进入对象类型要经过哪几站
  • 写出声明数据流和绑定对象类型的两条命令,并解释每个参数
  • 说出「只读」在每一站的具体含义,并选对主键

四站:从数据源到对象类型

前几门课画好了对象类型和关系,这一课让它们有数据。示例工厂有五个数据源:ERP(订单、物料、采购)、MES(工单、产出)、质量系统(检验)、设备物联网(读数)和 OA(审批)。它们都只读地进来:接入链路只从它们读取。Action 缺省只改语义层;配置 writeback webhook 后可以调用源系统。

  1. 01数据源只给一个只读账号
  2. 02发布端装在旁边,变了才推
  3. 03数据流有主键和字段白名单,按序落账
  4. 04同步绑定把流字段映射到属性
  5. 05对象类型workOrder、inspection……

发布端是你写的一个常驻小程序,外加 aidc semantic streams pipe:小程序用数据源自己的方式发现变化,pipe 负责算差、重试和心跳,只把变化的行发出去。发布端手里只有一把发布 Key,它只能往一条数据流发布。应用和智能体从不直连 ERP。

发布端怎样知道「变了」,取决于数据源:

源头推送旁边的增量探针
做法数据源在业务事件、流程节点或数据库变更时,主动调用发布端发布端在旁边问一个便宜的问题:版本号变了吗?更新时间的水位往前走了吗?变了才取全表
什么时候用数据源能推送,就一律用它数据源不能推送时的办法
对源系统的压力几乎没有一个探针服务所有订阅者,不是每个看板一个

发布端里那个小程序叫适配器,用什么语言都行,往标准输出每行打印一个 JSON。带 rows 的一行是有主键的流的全表,aidc semantic streams pipe 会和上一次比较,只发变化的行,第一次发全量;带 event 的一行是一条通知;带 status 的一行是适配器自报的健康状态,每一轮至少要有一行。

先声明数据流:写下主键和允许出现的字段。字段之外的内容会被拒收,证件号、银行卡、密码这类字段直接不收。

# 1. 声明:MES 工单流,主键是工单号
aidc semantic streams create mes-work-orders --title "MES 工单" --key work_order_no \
  --field work_order_no:string:工单号 --field status:string:状态 \
  --field planned_qty:number:计划数 --field good_qty:number:合格数 --field scrap_qty:number:报废数 \
  --field due_at:datetime:交期 --field line_code:string:产线

# 2. 签一把只能往这条流发布的 Key,装在发布端
aidc semantic streams key mes-work-orders --label "工厂服务器 · MES 发布端"

# 3. 在 MES 旁边常驻运行(systemd 服务,不是定时器)
AIDC_PUBLISH_KEY=aidc-pk-… aidc semantic streams pipe cell-example/mes-work-orders -- python3 mes_adapter.py

绑定:把流字段映射到属性

数据流只是原始信号。同步绑定写在 Object Type 定义的 datasources 里。它告诉本体(Ontology):这条流的每一行是一个 workOrder 对象,哪个流字段填哪个属性。对象类型要先定义好(aidc semantic define),绑定才建得起来。

aidc semantic datasource set workOrder --stream mes-work-orders \
  --map workOrderNo=work_order_no,status=status,plannedQty=planned_qty,goodQty=good_qty,scrapQty=scrap_qty,dueAt=due_at,lineCode=line_code
aidc semantic datasource list          # 每条绑定:状态、已同步到的流序号、落后几批
  • --map property=streamField:主键属性必须全部映射;--dry-run 只检查不落库。
  • --mode mirror(显式选择):流里的全部行就是全部对象,行从流里消失,对象就标成「源头已消失」。upsert(缺省):按主键增改,本批未出现的对象保留。
  • 绑定建好立即全量同步一次;之后流每落账一批,就在同一次发布里同步这一批的净变化。
  • 同步失败只记在绑定上(status: error 加原因),不影响数据流;下一批发现序号接不上会自动全量重同步,也可以手动 aidc semantic datasource resync workOrder --stream mes-work-orders。
  • 个别行映射不上(比如主键缺值)不会拦住整批:其余行照常同步,没能映射的行数和前几条原因记在绑定的 detail 里。

主键是绑定里最要紧的一列。它必须每行唯一,而且每次都一样:绑定之前,先在源数据里查一遍有没有重复的主键值。人改过的值和关系都挂在主键上,主键一变,就成了另一个对象,旧对象上的改动跟不过去。

随口一问

主键取导出表的行号,或者每次同步随机生成一个 id。

好的交代

主键取 MES 自己的工单号 workOrderNo:唯一,稳定,重导一遍也不变。

只读,是每一站的规矩

SOURCE

源系统

给发布端一个只读账号。能推送的源用推送;不能推送的,在旁边做一个便宜的增量探针,一个探针服务所有订阅者。

KEY

发布 Key

以 aidc-pk- 开头,只能读写它自己那一条流,读不了别的流,也进不了别的数据。

SEMANTIC

Semantic 这一侧

同步只写对象的数据源层。人和智能体的改动走 Action,落在另一层。配置 Action 的 writeback webhook 后,可以调用源系统。

这条链路的上限写在文档里,设计流的时候要留意:

1 MB单个事件上限
4 MB一条流最新状态上限
100每批最多条数
600每条流每分钟最多发布次数

下一节讲数据进来之后,人改过的值怎样和数据源的值合成对象的当前值。

要点

  • 数据经数据源、发布端、数据流、同步绑定进入对象类型,接入链路只读源系统。
  • 发布 Key 只能往一条流发布;应用和智能体从不直连数据源。
  • aidc semantic datasource set 把绑定写进 datasources,用 --map property=streamField 映射;主键必须全部映射,绑定建好立即全量同步一次。
  • 主键要唯一且稳定:用业务编号,不用行号或随机数。
  • Action 缺省只改语义层;配置 writeback webhook 后可调用源系统,调用失败则语义层不改。

练一练

接进示例工厂的采购订单

用蓝图里的 purchaseOrder(主键 poNo)练习:先纸上写,再在自己的公司里照着做一遍。

写出 aidc semantic streams create erp-purchase-orders …:主键 po_no,字段有状态、数量、预计到货日、供应商编号、物料编号,各写上类型。

小测

选一个答案,马上看解析。

Q1MES 导出表里,工单号是稳定的编号,行号每天重排。绑定时主键应该取?

Q2示例工厂的 MES 发布端拿到的发布 Key,能做什么?

Q3aidc semantic datasource set 成功之后,接下来会发生什么?

延伸阅读