Skip to content

传输配置

transports 段定义北向(northbound)传输通道,类比 Clash 的 proxy-groups。每个传输实例将采集到的数据点发布到云端或上游系统,并可选地接收反向写命令。

yaml
transports:
  - name: cloud-mqtt
    type: mqtt
    settings:
      broker: tcp://broker.emqx.io:1883
      client-id: factory-edge-01
      qos: 1
      topic-template: "factory/{{.Driver}}/{{.Group}}/{{.Tag}}"
      command-topic: "factory/commands/#"
    batch-size: 50
    flush-interval: 1s

传输通用字段

所有传输实例共享以下字段,settings 内容随 type 不同而变化。

字段类型必填默认值说明
namestring传输实例名称,全局唯一,用于规则 target 引用
typestring传输类型:mqtthttp
settingsobject协议专属连接参数
batch-sizeint100单次批量发送的最大数据点数
flush-intervalduration触发刷新的时间间隔,未设置则不启用定时刷新
retry-countint0(不重试)发送失败后的重试次数
buffer-sizeint100命令/数据通道容量,超出后丢弃

支持的传输类型:

type协议典型目标
mqttMQTT 3.x/5.0IoT 平台、边缘网关、命令回写
httpHTTP/RESTMES、数据湖 API、Webhook

INFO

batch-sizeflush-interval 共同决定发送时机:任一条件先满足即触发发送。例如 batch-size: 50flush-interval: 1s,则 1 秒内攒满 50 点立即发送,否则满 1 秒发送当前已攒的点。


通用参数详解

batch-size

单次批量发送的数据点上限。较大的值提升吞吐但增加单次延迟,较小的值降低延迟但增加请求次数。

场景推荐值
高频传感器、云端聚合100500
低频关键数据、需低延迟110
MQTT 受限于单包大小50100

flush-interval

定时刷新间隔,保证数据点不会因 batch-size 未凑满而无限滞留。采用 Go duration 字符串(1s500ms 等)。未设置时不启用定时刷新,仅按 batch-size 触发发送。

TIP

若数据稀疏且未设置 flush-interval,点会长期滞留队列,不建议用于实时场景。建议显式设置一个合理的刷新间隔。

retry-count

发送失败后的重试次数。每次重试间隔采用指数退避(base × 2^n)。默认 0 表示不重试,重试耗尽后数据点直接丢弃并累加 failed 计数。

离线缓冲

全局 buffer 已废弃,不再支持本地落盘缓冲。如需重试与批量,请使用 transport 级别的 retry-countbatch-sizeflush-interval

buffer-size

命令/数据通道容量(内部 command/data channel)。当上游采集速率瞬时超过传输速率,队列积压达到上限时,丢弃最旧的数据点并记录 warn 日志。默认 100


MQTT (mqtt)

通过 MQTT 协议发布数据点,并可选订阅命令主题接收反向写命令(下行控制)。

yaml
- name: cloud-mqtt
  type: mqtt
  settings:
    broker: tcp://broker.emqx.io:1883
    client-id: factory-edge-01
    qos: 1
    topic-template: "factory/{{.Driver}}/{{.Group}}/{{.Tag}}"
    command-topic: "factory/commands/#"
  batch-size: 50
  flush-interval: 1s

settings 字段

字段类型必填默认值说明
brokerstringBroker 地址,格式 scheme://host:port
client-idstring自动生成客户端 ID,需在 Broker 范围内唯一
usernamestring用户名认证
passwordstring密码认证
qosint1服务质量等级,取值 012
retainedboolfalse是否发布保留消息
topic-templatestringcorec/{{.Driver}}/{{.Tag}}发布主题模板,支持 Go 模板变量
command-topicstring订阅的命令主题,支持通配符 + / #
data-topicstring订阅的数据主题(链式内核入站),需配合 parser
keep-aliveduration60s心跳保活间隔
connect-timeoutduration10s连接超时时间
auto-reconnectbooltrue是否自动重连
clean-sessionbooltrue是否使用干净会话
connect-retrybooltrue是否在连接失败后自动重试
connect-retry-intervalduration5s连接重试间隔
subscribe-timeoutduration5s订阅操作超时
publish-timeoutduration5s发布操作超时
disconnect-quiesceduration1s断开连接时的静默等待时长
parserobject入站数据解析器配置,见解析器配置

broker

scheme含义示例
tcp://明文 MQTTtcp://broker.emqx.io:1883
ssl://TLS 加密 MQTTssl://broker.emqx.io:8883
ws://WebSocketws://broker.emqx.io:8083/mqtt
wss://WebSocket over TLSwss://broker.emqx.io:8084/mqtt

qos

取值含义适用场景
0至多一次,不保证送达高频低价值遥测
1至少一次,保证送达但可能重复推荐,工业数据常用
2恰好一次,开销最大命令回写、计费数据

WARNING

client-id 必须唯一。重复 ID 会导致后连接的客户端踢掉先前的会话,造成数据中断。多实例部署时建议附加主机名或 Pod 名后缀。

topic-template

发布主题使用 Go template 语法,可引用数据点字段:

变量来源示例值
{{.Driver}}驱动实例名plc-modbus
{{.Group}}标签分组sensors
{{.Tag}}标签名temperature
{{.Device}}设备名plc-01
yaml
topic-template: "factory/{{.Driver}}/{{.Group}}/{{.Tag}}"
# 渲染结果:factory/plc-modbus/sensors/temperature

TIP

主题层级设计建议遵循 组织/产线/设备/测点 的从宽到窄结构,便于下游按层级订阅与权限控制。

command-topic

订阅该主题以接收反向写命令。消息体为 JSON 格式的 WriteCommand

json
{
  "driver": "plc-modbus",
  "device": "plc-01",
  "tag": "pump_status",
  "value": true,
  "type": "bool"
}

支持 MQTT 通配符:

通配符含义
+单层通配,匹配一个主题层级
#多层通配,匹配剩余所有层级(只能放末尾)
yaml
command-topic: "factory/commands/#"   # 接收所有命令
command-topic: "factory/commands/plc-modbus/+"   # 仅该驱动的命令

data-topic(链式内核入站)

订阅数据主题以接收上游 CoreC 实例或第三方 MQTT 发布者的数据(链式内核 / chained-kernel 入站)。设置 data-topic 后需配合 parser 配置将消息体解析为 DataPoint。解析后的数据点进入引擎 DataBus,流经规则并可被重新发布到下游传输,实现多跳中继拓扑(边缘 → 网关 → 云端)。

yaml
settings:
  data-topic: "upstream/factory/#"
  parser:
    type: default          # json.Unmarshal(DataPoint),CoreC→CoreC 零成本路径

HTTP (http)

通过 HTTP/REST 将数据点批量推送到上游 API,适用于 MES、数据湖、Webhook 等场景。HTTP 传输也可通过 webhook 接收入站数据(链式内核入站)。

WARNING

HTTP Push 的反向控制指令通道(OnCommand())存在但当前不投递。反向下发请使用 MQTT command-topic

yaml
- name: mes-http-push
  type: http
  settings:
    url: "http://localhost:8080/api/v1/telemetry"
    method: POST
    headers:
      Authorization: "Bearer mes-secret-key"
      Content-Type: "application/json"
    timeout: 3s
  batch-size: 100
  flush-interval: 5s

settings 字段

字段类型必填默认值说明
urlstring上游 API 地址
methodstringPOSTHTTP 方法,通常 POSTPUT
headersmap[string]string自定义请求头
timeoutduration5sHTTP 请求超时时间
max-idle-connsint100连接池最大空闲连接数
max-idle-conns-per-hostint20每主机最大空闲连接数
idle-conn-timeoutduration90s空闲连接超时时间
webhook-addrstringWebhook 监听地址(链式内核入站),如 0.0.0.0:9091
webhook-pathstring/dataWebhook URL 路径
parserobject入站数据解析器配置,见解析器配置

url

完整 URL,包含 scheme 与路径。支持 http://https://,使用 https:// 时内核自动校验服务端证书。

yaml
url: "https://mes.example.com/api/v1/telemetry"

method

取值含义
POST创建新记录,最常用
PUT整体替换资源
PATCH部分更新

headers

键值对形式的自定义请求头,常用于鉴权与内容类型声明:

yaml
headers:
  Authorization: "Bearer mes-secret-key"
  Content-Type: "application/json"
  X-Tenant-Id: "factory-01"

INFO

内核默认以 JSON 数组形式发送批量数据点,即 []DataPoint。若上游要求不同结构,可在网关层做转换,或通过规则 transform 调整。

timeout

单次 HTTP 请求(含重试中的每次尝试)的超时时间,默认 5s。超时后视为本次发送失败,进入重试流程。

连接池

以下字段控制底层 HTTP 客户端的连接池行为,均为可选:

字段默认值说明
max-idle-conns100连接池最大空闲连接数
max-idle-conns-per-host20每主机最大空闲连接数
idle-conn-timeout90s空闲连接超时时间

webhook-addr / webhook-path(链式内核入站)

设置 webhook-addr 后,传输会启动一个 HTTP 服务接收上游 CoreC 实例或第三方 POST 请求中的数据(链式内核入站)。请求体可以是单个 DataPoint JSON 对象或 []DataPoint 数组,经 parser 解析后进入引擎 DataBus。

yaml
settings:
  webhook-addr: "0.0.0.0:9091"
  webhook-path: "/ingest"
  parser:
    type: default          # json.Unmarshal(DataPoint)

解析器配置(parser)

当传输配置了入站通道(MQTT data-topic 或 HTTP webhook-addr)时,需通过 parser 配置将消息体解析为 DataPointparsersettings 下的一个子对象。

parser.type

取值说明
default直接 json.UnmarshalDataPoint,CoreC→CoreC 零成本路径(默认)
jsonpath通过 Go template 将任意 JSON 字段映射到 DataPoint
raw将整个 payload 视为标量值,tag 可从 MQTT topic 提取

default 解析器

无需额外配置,消息体须为 DataPoint 的 JSON 表示。

jsonpath 解析器

字段类型必填默认值说明
parser.driverstring静态 driver 字段值
parser.tagstringtag 字段的模板或静态字符串
parser.valuestringvalue 字段的模板路径
parser.data-typestring数据类型名,如 float32
parser.groupstring静态 group 字段值
parser.devicestring静态 device 字段值
parser.timestampstringtimestamp 字段的模板路径
parser.timestamp-formatstringrfc3339时间戳格式:rfc3339unixunixmilli
yaml
parser:
  type: jsonpath
  driver: "lora-gateway"
  tag: "{{ .payload.dev_id }}"
  value: "{{ .payload.temp }}"
  data-type: "float32"
  group: "sensors"
  timestamp: "{{ .payload.ts }}"
  timestamp-format: "unix"

raw 解析器

将整个 payload 视为标量值,tag 可从 MQTT topic 的指定层级提取。

字段类型必填默认值说明
parser.driverstring静态 driver 字段值
parser.tagstring静态 tag 字段值
parser.tag-from-topicint0取 MQTT topic 第 N 段作为 tag,0 表示使用静态 tag
parser.data-typestringfloat64数据类型名
parser.groupstring静态 group 字段值
parser.devicestring静态 device 字段值
parser.timestamp-formatstringrfc3339时间戳格式
yaml
parser:
  type: raw
  driver: "factory"
  tag-from-topic: 2      # 使用 MQTT topic 第 2 段作为 tag
  data-type: "float32"

完整示例

yaml
transports:
  # MQTT 云端传输:发布 + 命令回写 + 链式入站
  - name: cloud-mqtt
    type: mqtt
    settings:
      broker: ssl://broker.emqx.io:8883
      client-id: factory-edge-01
      qos: 1
      topic-template: "factory/{{.Driver}}/{{.Group}}/{{.Tag}}"
      command-topic: "factory/commands/#"
      data-topic: "upstream/factory/#"
      parser:
        type: default
    batch-size: 50
    flush-interval: 1s
    retry-count: 5
    buffer-size: 100

  # HTTP 推送:MES 数据接入 + Webhook 入站
  - name: mes-http-push
    type: http
    settings:
      url: "https://mes.example.com/api/v1/telemetry"
      method: POST
      headers:
        Authorization: "Bearer mes-secret-key"
      timeout: 5s
      webhook-addr: "0.0.0.0:9091"
      webhook-path: "/ingest"
      parser:
        type: default
    batch-size: 100
    flush-interval: 5s
    retry-count: 3
    buffer-size: 100

TIP

对关键数据建议同时配置 MQTT 与 HTTP 两个传输,并通过规则 mirror 动作双发,实现传输层冗余。

Released under the MIT License.