Skip to content

链式内核(Chained Kernel)

链式内核让 CoreC 实例既能接收上游数据,又能转发到下游。一个 Transport 配了 data-topic(MQTT)或 webhook-addr(HTTP)就开启接收,收到的数据经 Parser 解析后喂回 DataBus,走完规则匹配再发布到下游——实现 edge → gateway → cloud 的多跳中继。

灵感来自 mihomo 的 relay 架构:inner listener 把连接回流到 tunnel,CoreC 用 channel 代替管道,把结构化的 DataPoint 回流到 processingLoop。

数据流

Driver.Read()  ──▶ onDriverData() ──▶ DataBus.Push() ─┐
                                                        ├──▶ processingLoop ──▶ 规则 ──▶ Publish
Transport.OnData() ──▶ startDataListener() ──▶ DataBus.Push() ─┘

两个数据来源(驱动轮询、Transport 接收)汇入同一条 DataBus,对 processingLoop 完全透明。

判断一个节点是"采集"还是"中继"

采集节点中继节点
drivers有驱动配置drivers: []
Transport只有 topic-template / url额外配了 data-topic / webhook-addr
数据来源Driver.Read()Transport.OnData()

Parser 三种模式

模式适用场景原理
defaultCoreC → CoreCjson.Unmarshal(DataPoint),格式天然匹配,零成本
jsonpath第三方 JSONtext/template 字段映射,双花括号模板语法
raw裸数值payload 即 value,tag 从 topic 路径提取

场景一:纯订阅转发(协议网关)

CoreC 不连任何设备,只订阅第三方设备的 MQTT,解析后转发到云端。角色是协议适配器

LoRa网关 ──MQTT──▶ CoreC ──MQTT──▶ 云端
(非CoreC格式)       (解析转换)      (标准格式)
yaml
drivers: []

transports:
  - name: bridge
    type: mqtt
    settings:
      broker: tcp://broker:1883
      topic-template: "cloud/{{.Driver}}/{{.Tag}}"
      data-topic: "lora/+/up"              # 订阅第三方
      parser:
        type: jsonpath                      # 第三方格式 → DataPoint
        driver: "lora"
        tag: "{{ .payload.dev_id }}"
        value: "{{ .payload.temp }}"
        data-type: "float32"

rules:
  - { match: "ALL", action: forward, target: bridge }

特点drivers: [],数据全从 data-topic 进来。CoreC 只做格式转换。


场景二:两个 CoreC 级联(边缘 → 云端)

最常见的链式拓扑。A 在边缘连 PLC 采集,B 在云端收数据做告警 / 存储。

PLC ──modbus──▶ CoreC-A ──MQTT──▶ CoreC-B ──MQTT──▶ 云端
                 (边缘)            (中继)

CoreC-A(边缘):

yaml
drivers:
  - name: plc
    type: modbus-tcp
    settings: { host: 192.168.1.100, port: 502 }
    tags:
      - { name: temperature, address: "40001", type: float32, interval: 1s }

transports:
  - name: mqtt-out
    type: mqtt
    settings:
      broker: tcp://broker:1883
      topic-template: "edgeA/{{.Driver}}/{{.Tag}}"
      # 没有 data-topic,A 不收数据

rules:
  - { match: "ALL", action: forward, target: mqtt-out }

CoreC-B(中继):

yaml
drivers: []

transports:
  - name: mqtt-chain
    type: mqtt
    settings:
      broker: tcp://broker:1883
      topic-template: "cloud/{{.Driver}}/{{.Tag}}"
      data-topic: "edgeA/#"               # 收 A 发的
      parser:
        type: default                      # A 发的就是 DataPoint JSON

rules:
  - { match: "tag == 'temperature' && value > 90", action: alert, target: mqtt-chain, priority: 1 }
  - { match: "ALL", action: forward, target: mqtt-chain, priority: 999 }

特点:A 的 topic-template 和 B 的 data-topic 必须匹配。A 发 edgeA/plc/temperature,B 订阅 edgeA/# 就能收到。


场景三:多级级联(边缘 → 网关 → 云端)

三跳或更多。每一跳都可以做规则处理。

PLC ──▶ CoreC-A ──MQTT──▶ CoreC-B ──HTTP──▶ CoreC-C ──MQTT──▶ 云端
         (车间)            (厂区网关)          (区域中心)

CoreC-B(厂区网关,MQTT 收 → HTTP 发):

yaml
drivers: []

transports:
  - name: from-edge
    type: mqtt
    settings:
      broker: tcp://factory-broker:1883
      topic-template: "unused"
      data-topic: "edgeA/#"
      parser: { type: default }

  - name: to-region
    type: http
    settings:
      url: "http://corec-c:9091/ingest"    # 发给 C 的 webhook

rules:
  - { match: "ALL", action: forward, target: to-region }

CoreC-C(区域中心,HTTP 收 → MQTT 发):

yaml
drivers: []

transports:
  - name: from-gateway
    type: http
    settings:
      url: "http://localhost:9999/unused"
      webhook-addr: "0.0.0.0:9091"
      webhook-path: "/ingest"
      parser: { type: default }

  - name: to-cloud
    type: mqtt
    settings:
      broker: tcp://cloud:1883
      topic-template: "cloud/{{.Driver}}/{{.Tag}}"

rules:
  - { match: "ALL", action: forward, target: to-cloud }

特点:中间每一跳都可以做过滤、告警、变换。比如 B 过滤掉车间不关心的数据,C 做区域级告警。


场景四:协议转换(MQTT → HTTP)

CoreC 收 MQTT 数据,通过 HTTP 推到 REST API。解决云端只接受 HTTP,设备只发 MQTT 的情况。

设备 ──MQTT──▶ CoreC ──HTTP──▶ MES/ERP系统
yaml
drivers: []

transports:
  - name: mqtt-in
    type: mqtt
    settings:
      broker: tcp://broker:1883
      topic-template: "unused"
      data-topic: "factory/#"
      parser: { type: default }

  - name: http-out
    type: http
    settings:
      url: "http://mes.example.com/api/v1/telemetry"
      headers:
        Authorization: "Bearer key"

rules:
  - { match: "ALL", action: forward, target: http-out }

反过来也行(HTTP → MQTT):配 webhook-addr 收,topic-template 发。


场景五:多对一汇聚

多个边缘节点往同一个中继汇聚,中继统一处理后转发。

CoreC-A ──MQTT──▶
CoreC-B ──MQTT──▶ CoreC-G ──MQTT──▶ 云端
CoreC-C ──MQTT──▶
yaml
# CoreC-G(汇聚网关)
drivers: []

transports:
  - name: aggregator
    type: mqtt
    settings:
      broker: tcp://broker:1883
      topic-template: "cloud/{{.Driver}}/{{.Tag}}"
      data-topic: "edge/#"            # # 匹配所有子层,收 edge/A/... edge/B/... edge/C/... 所有节点
      parser: { type: default }

rules:
  - { match: "ALL", action: forward, target: aggregator }

A 发 edge/A/...,B 发 edge/B/...,C 发 edge/C/...,G 订阅 edge/# 全收。

MQTT 单层通配符 + 必须独占一整个层级(如 edge/+/up),不能写成 edge+/#。若只需匹配一层节点名,可用 edge/+/... 形式并让各节点发布到 edge/A/...edge/B/...


场景六:一对多分发

一个边缘节点采集数据,同时推到多个云端(MQTT + HTTP,或两个不同 MQTT broker)。

              ┌──MQTT──▶ 云端A
PLC ──▶ CoreC ┤
              └──HTTP──▶ MES系统
yaml
drivers:
  - name: plc
    type: modbus-tcp
    settings: { host: 192.168.1.100, port: 502 }
    tags:
      - { name: temperature, address: "40001", type: float32, interval: 1s }

transports:
  - name: cloud-mqtt
    type: mqtt
    settings:
      broker: tcp://cloud-a:1883
      topic-template: "cloud/{{.Driver}}/{{.Tag}}"

  - name: mes-http
    type: http
    settings:
      url: "http://mes.example.com/api/telemetry"

rules:
  - { match: "ALL", action: mirror, targets: [cloud-mqtt, mes-http] }

特点:用 action: mirror 同时推多个 target。这不是链式内核的新功能,但可以跟链式组合。


场景七:双向级联(数据上行 + 命令下行)

数据从边缘往云端走,控制命令从云端往边缘走。两个方向都经过中继。

         数据上行                    数据上行
PLC ◀──modbus──▶ CoreC-A ◀──MQTT──▶ CoreC-B ◀──MQTT──▶ 云端
         写命令                    写命令
         下行                       下行

CoreC-A(边缘):

yaml
transports:
  - name: mqtt
    type: mqtt
    settings:
      broker: tcp://broker:1883
      topic-template: "edgeA/{{.Driver}}/{{.Tag}}"
      command-topic: "edgeA/commands/#"    # 收 B 转发下来的命令

CoreC-B(中继):

yaml
transports:
  - name: mqtt-chain
    type: mqtt
    settings:
      broker: tcp://broker:1883
      topic-template: "cloud/{{.Driver}}/{{.Tag}}"
      data-topic: "edgeA/#"                # 收 A 的数据
      command-topic: "cloud/commands/#"    # 收云端的命令
      parser: { type: default }

命令跨实例转发当前不支持

上图为目标拓扑,但 CoreC 当前没有命令重发布机制:当 B 通过 command-topic 收到云端命令时,它只会调用自身本地驱动Write()(而 B 是中继节点,drivers: [],没有本地驱动可写),不会把命令重新发布到 edgeA/commands/#。因此命令无法从 B 透传到 A。

目前可行的替代方案:

  • 方案一:云端直接向 A 的 command-topicedgeA/commands/#)下发命令,跳过中继 B。
  • 方案二:在中继 B 上用外部脚本 / 规则引擎监听 cloud/commands/# 并转发到 edgeA/commands/#
  • 方案三:等待后续版本支持 command 透传 / 重发布能力。

每个 CoreC 实例的 command-topic 只触发自身本地驱动的写入,不会跨实例传递命令。


场景总览

场景拓扑CoreC 角色链式内核关键配置
① 纯订阅转发第三方 → CoreC → 云端协议适配器data-topic + parser: jsonpath
② 两级级联CoreC-A → CoreC-BA=采集,B=中继B 用 data-topic + parser: default
③ 多级级联A → B → C每跳都可处理中间节点 data-topic / webhook-addr
④ 协议转换MQTT → CoreC → HTTP协议桥data-topic 收,url
⑤ 多对一汇聚A+B+C → G → 云端汇聚网关G 用 data-topic: "edge/#" 收所有 edge/A/...
⑥ 一对多分发CoreC → 多个云端边缘采集action: mirror + 多 target
⑦ 双向级联数据上行 + 命令下行中继双向data-topic + command-topic

配置速查

MQTT Transport

参数作用必填
brokerMQTT broker 地址
topic-template发数据的 topic 模板有默认值
command-topic收写命令的 topic可选
data-topic收数据的 topic(开启链式内核)可选
parser数据解析器(配了 data-topic 就必须配)条件必填

HTTP Transport

参数作用必填
url发数据的目标 URL
webhook-addr收数据的监听地址(开启链式内核)可选
webhook-path收数据的 URL 路径(默认 /data可选
parser数据解析器(配了 webhook-addr 就必须配)条件必填

Parser 配置

yaml
# CoreC → CoreC(零成本,格式天然匹配)
parser:
  type: default

# 第三方 JSON(字段名不同)
parser:
  type: jsonpath
  driver: "lora-gateway"
  tag: "{{ .payload.dev_id }}"
  value: "{{ .payload.temp }}"
  data-type: "float32"
  group: "sensors"
  timestamp: "{{ .payload.ts }}"
  timestamp-format: "unix"          # rfc3339 | unix | unixmilli

# 裸数值(payload 即 value)
parser:
  type: raw
  driver: "factory"
  tag: "temperature"                # 静态 tag
  # 或 tag-from-topic: 2            # 从 topic 第 2 段取 tag
  data-type: "float64"

Released under the MIT License.