连接 SDK
把应用接到 AIDC 的其他部分:企业数据集(固定报告)、实时数据流与智能体。
import { connect } from "/developer/sdk/v1/aidc.js";
数据集
const list = await connect.datasets.list(); // 调用方能读的数据集
const report = await connect.datasets.report("aidc-demo-erp", "overview");
平台连接企业数据源,只出固定口径的聚合报告,不出明细行。应用只能读清单 datasets 里声明的数据集。
当前可用:
| 数据集 | 说明 |
|---|---|
aidc-demo-erp |
AIDC 合成 ERP 演示库(AIDC 数据库里的 demo_erp schema):虚构企业 12 个月的客户、供应商、物料、销售 / 采购订单、应收应付、收付款、库存流水。全部为虚构数据,给公开样板用 |
overview 报告包含:营收与毛利、订单与客单价、应收与逾期、DSO、库存金额、30 天内应付;月度趋势;客户 Top 10(含上下半年走势);品类结构;应收账龄;逾期客户;低于安全库存与呆滞物料;供应商采购额;以及规则化的经营信号(大客户下滑、逾期集中、营收异动、库存风险、回款周期)。
客户自己的 ERP(金蝶、用友、SAP、SQL Server…):连接在客户侧执行(数据不出客户环境),平台登记数据源、出同样形状的报告——接入由 AIDC 交付团队配置。
实时数据流
数据在源头(ERP、MES、SAP、OA…)一变,就由源头旁边的发布端推给 AIDC,AIDC 按序落账、保存最新状态,再在几秒内推给所有订阅端(Nexus 应用、CLI、智能体)。应用不轮询、不跑 cron,也不会有十几个看板各自去查同一个 ERP。
数据源 ──(发布端:aidc stream pipe + 适配器,装在数据源旁边)──HTTPS──▶ AIDC(按序落账 + 最新状态)──SSE──▶ 订阅端
订阅(在应用里)
清单里声明要订阅的流(与应用同一命名空间):"streams": ["orders"]。
const orders = connect.stream("orders");
const sub = orders.subscribe({
onUpdate({ kind, rows, upserted, deleted, sourceAt }) {
render(rows); // 合并好的当前全部行(按主键)
flash(upserted); // 这次变化的行
},
onNotice({ data }) { toast(data.message); }, // 发布端发的通知(如「新合同」)
onStatus({ connection, publisher }) {
// connection: connecting / live / reconnecting / closed
// publisher.live = false → 提示「数据可能不是最新」;publisher.status = "source_unreachable" → 「连不上数据源」
badge(connection, publisher);
},
});
const { rows, seq } = await orders.get(); // 一次性读当前状态
- 第一次连上先给全量,之后每次变化给增量;断线自动重连并带着最后的序号续传——不丢、不重。
- 每条流都带发布端的健康状态:在线 / 离线(90 秒没心跳)、正常 / 部分异常 / 连不上数据源 / 已停止。数据源断了,订阅端看得见,不会像 cron 看板那样静默变旧。
声明与发布
# 1. 声明:主键 + 字段白名单(发布的数据只许出现这些字段;证件号、银行卡、密码等类别直接拒绝)
aidc stream create orders --title "在手订单" --key order_no \
--field order_no:string:合同号 --field qty:number:数量 --field due:datetime:交期
# 2. 签一把发布 Key:只能往这一条流发布,装在发布端(.env 0600)
aidc stream key orders --label "工厂服务器 · 订单发布端"
# 3. 在数据源旁边常驻运行发布端(systemd 服务,不是定时器)
AIDC_PUBLISH_KEY=aidc-pk-… aidc stream pipe cell-xxx/orders -- python3 orders_adapter.py
适配器是你写的常驻小程序(任何语言),用数据源自己的方式发现变化,往 stdout 每行打印一个 JSON:
| 行 | 含义 |
|---|---|
{"rows": [...]} |
有主键的流的全表——发布端和上一次比较,只发变化的行(第一次发全量) |
{"value": {...}} |
没有主键的流的整体值——变了才发 |
{"event": {...}} |
通知(字段同样受白名单约束,另可带 message) |
{"status": "ok" | "degraded" | "source_unreachable", "detail": "…"} |
适配器自报健康;每轮至少打一行,发布端 15 分钟收不到任何输出就报 degraded |
aidc stream pipe 负责其余的事:算差、带幂等 id 重试(发布失败不丢、不重)、每 30 秒心跳、每 6 小时发一次全量校准、适配器崩溃按退避重启并上报「连不上数据源」。
# orders_adapter.py —— 常驻;连接保持;先用便宜的探针看有没有变化,变了才取全表
import json, time, pymssql
conn, last = None, None
while True:
try:
conn = conn or pymssql.connect(server="…", user="readonly", password=open("/path/secret").read().strip(), database="erp")
cur = conn.cursor(as_dict=True)
cur.execute("SELECT CHANGE_TRACKING_CURRENT_VERSION() AS v") # 没开 Change Tracking 就用 rowversion / 更新时间水位
version = cur.fetchone()["v"]
if version != last:
cur.execute("SELECT order_no, qty, due FROM dbo.orders WHERE status = 'open'")
print(json.dumps({"rows": cur.fetchall()}, default=str), flush=True)
last = version
else:
print('{"status": "ok"}', flush=True)
except Exception as err:
print(json.dumps({"status": "source_unreachable", "detail": str(err)[:200]}), flush=True)
conn = None
time.sleep(5)
也可以不用 CLI:在 Node / 服务端用 connect.stream("cell-xxx/orders").publisher()(rows() 自动算差、event()、status()),或直接 POST /api/v1/developer/streams/{namespace}/{name}/events。
怎么「发现变化」:按数据源
发布端之后的链路(落账、推送、续传、鉴权)对所有数据源都一样;不同的只是适配器怎么知道「变了」。能推送的源一律用推送;源头不能推送时,才在源头旁边做便宜的增量探针——一个探针服务所有订阅者,不是每个看板一个 cron。
| 数据源 | 最佳:源头推送 | 次选:边缘增量探针 |
|---|---|---|
| SAP ECC / S/4HANA | S/4 业务事件(Event Mesh)、IDoc / 变更指针、BAdI 调 HTTP | RFC 按录入时间 / 计数器增量读(例:报工确认 AFRU 的 ERSDA/ERZET) |
| SQL Server(自研 / 行业 ERP) | Change Tracking / CDC(DBA 开启一次) | rowversion / 更新时间列水位;都没有时对「当天范围」做汇总指纹 |
| PostgreSQL(MES 等) | 主库逻辑复制 / LISTEN-NOTIFY | 只读热备:WAL 回放位置门控 + 限频查询 |
| 金蝶 / 用友 / Odoo / SaaS ERP | 开放平台 webhook、业务事件订阅 | 开放 API 按修改时间增量拉 |
| OA(泛微等) | 流程节点动作调 HTTP | 流程表 requestid / 修改时间水位 |
| 企业微信 / 钉钉 | 事件回调 / Stream 模式 | — |
| 邮件 | IMAP IDLE | — |
为什么不需要 Kafka:一条流的量级是每秒个位数事件、几个到几十个订阅者;要的是「有序、断线可补齐、新订阅者先拿最新状态、按公司鉴权」。AIDC 用数据库账本(有序、可补齐)+ 最新状态 + 已有的低延迟中继(唤醒订阅端)就够了,发布端只需要能出 HTTPS。到了每秒成千上万条事件、需要长时间重放给很多独立消费方时,再引入日志型中间件;即便那时,Debezium 这类 CDC 也可以直接把变化 POST 给同一个发布接口。
限制与权限
- 单个事件 ≤ 1 MB、最新状态 ≤ 4 MB、每批 ≤ 100 条、每条流每分钟 ≤ 600 次发布;默认保留最近 1000 条事件供断线补齐(更早断开的订阅者直接拿最新状态)。
- 读:清单声明了这条流的同命名空间应用(客户应用只签给已登录的本公司成员)、本公司的开发者 Key;发布 Key 只能读、写它自己那一条流。
- 订阅连接约每 4 分 40 秒由服务端收尾一次,SDK 自动续接,应用无感。
- 删除数据流(
aidc stream delete)会物理删除它的事件账与状态,并吊销它的发布 Key。
API
| 方法 | 路径 | 说明 |
|---|---|---|
GET / POST |
/api/v1/developer/streams |
列出 / 声明(幂等) |
GET / DELETE |
/api/v1/developer/streams/{namespace}/{name} |
信息 + 最新状态 / 删除 |
GET |
/api/v1/developer/streams/{namespace}/{name}/events |
订阅:Accept: text/event-stream 返回 SSE(state / change / notice / publisher / reconnect,每条 id = 序号,?after=N 续传);否则返回 N 之后的 JSON 增量 |
POST |
/api/v1/developer/streams/{namespace}/{name}/events |
发布 {events:[{type, id?, at?, data}]}(type = snapshot / patch / event / status) |
GET / POST / DELETE |
/api/v1/developer/streams/{namespace}/{name}/keys |
发布 Key:列出 / 签发 / 吊销(?prefix=) |
智能体
智能体走现成的 OpenAI 兼容 Agent API——与 Adis、外部框架调智能体是同一条路。凭证是访客自己的 Agent Key(aidc-sk-…,在 Console 的智能体页或 Adis 里获取);SDK 只在内存里使用,不存储、不上报。
const agent = connect.agent(agentKey);
const [info] = await agent.info(); // 这把 Key 能调用的智能体
const reply = await agent.send("请把这份会议纪要归档,并跟进待办。\n\n" + minutes.markdown, { sessionId: "meeting-0925" });
for await (const delta of agent.stream("今天 3 号线有几件不合格?")) out.textContent += delta;
const fileId = await agent.upload(photo.blob, "defect.jpg"); // 图片 / 文件,24 小时内有效
await agent.send("请登记这件不合格品。", { files: [fileId] });
sessionId相同即同一段会话,智能体记得上文。- 文件上传只支持云端运行的智能体;本地节点 / 快速车道的 Key 会报错。
- 在 AIDC 站点以外(CLI、服务端)用:
connect.agent(key, { agentBase: "https://api.ai-dc.ai/v1" })。
命令行
aidc connect datasets
aidc connect report aidc-demo-erp --json
AIDC_AGENT_KEY=aidc-sk-… aidc connect send "你好" --session demo
aidc stream create <流名> --title … --key <主键> --field 名:类型[:显示名] …
aidc stream key <流名> --label "装在哪" # 发布 Key(aidc-pk-…),只显示一次
AIDC_PUBLISH_KEY=aidc-pk-… aidc stream pipe <命名空间>/<流名> -- <适配器命令…>
aidc stream tail <流名> # 实时订阅,每次变化一行
aidc stream status <流名> | get <流名> | list | keys | revoke | publish | delete
本页由 developer/docs/connect.md 生成 · Markdown 原文 · llms.txt