给属性上报链路做基线是一件很枯燥的事:固定 10k 设备在线、10k msg/s、每条报文 500 字节上下、平均一条报文带 5 个属性,然后让 benchmark 跑三轮。跑完拿到的第一个数字有点刺眼——一次属性上报,在热路径上要做67 次堆分配。
67 次里面真正属于业务的没几次。剩下的大部分花在同一个值被反复改名、反复转换上:进来时叫报文字节,解完叫 params,归一后叫 ReportPropertyData,发给下游时又要变成各自投影过的形状。
这篇文章讲的是这条链路上做过的一次收敛:把"同一个值"在平台内部的表述固定下来,让影子、告警、时序库、审计日志四个消费者吃同一份数据;顺带把分配次数从 67 降到 37。为了让说法有据可查,文中凡是"代码里是这么做的",指的都是本地发行版(Professional)的源码——开源主干里不包含 Poller 与驱动插件这套东西。
一、一条属性上报有三条来路
SagooIoT 的设备接入有三条主路径,属性上报从任意一条都能进来。
| 路径 | 入口 | 协议解码 | JS 脚本 |
|---|---|---|---|
| MQTT 直连设备 | mqtt.Subscribe→HandleMessage | 走协议插件 | 走 |
| 主动采集设备(Poller) | Poller Task→ 驱动插件 Read 或 DTU 隧道 Ask | 跳过(采集侧已是 JSON) | 走 |
| 通道接入设备(Tunnel) | Server Receive→TunnelBase.ReadData | 必须走(要靠解码结果路由) | 走 |
三条路最后都汇到同一个函数:
funcHandleDirectMessage(ctx context.Context,topicstring,payload[]byte,handlerfunc(context.Context,topicModel.TopicHandlerData)error,skipProtocolDecodebool,logTypestring)error{...res,err:=processMessage(ctx,processDevice,payload,logType,skipProtocolDecode)...}Poller 和 Tunnel 调它时skipProtocolDecode传 true,因为它们已经各自完成了协议解码——一个在采集侧、一个在 Tunnel 层。
这套"多源汇聚、统一管线"不是一开始就这样。Tunnel 路径原来有一份独立的router()实现,后来被删除合流,原因是它带着三个具体问题:不支持 JS 的reply函数、不受全局协程池(gPool)背压控制、不被 before/after Filter 拦截。三条都不是理论缺陷,是能在现场撞到的行为差异。
二、同一个值,四个名字
收敛的前提是先把"值"在链路上的形态命名清楚。整条链路里同一个业务值最多有四种形态:
| 代号 | 类型 | 含义 |
|---|---|---|
| Raw | owned[]byte/string | 协议解码 + JS 脚本处理之后的报文原文 |
| A | map[string]any | 反序列化后的协议树(params) |
| B / Canonical | ReportPropertyData | 过物模型过滤、做过类型与时间归一之后的"平台认定值" |
| Derived envelope | 平台重编的 JSON | 不是原文,只用于没有独立原文的场景 |
Raw 由processMessage产出。它返回的不是裸[]byte,而是一个包装类型:
// network/core/owned_payload.go// OwnedPayload 持有独立分配的 payload 字节,调用方可安全跨异步持有。typeOwnedPayloadstruct{data[]byte}Canonical 由effects.ParseProperties产出——它做三件事:按产品物模型挑出有定义的属性、按 TSL 声明把值转成平台认定的类型、补时间戳。Derived envelope 则出现在网关批量里:子设备的属性嵌在网关的报文里,没有自己的报文原文,平台只能自己拼一份 JSON 出来。
四种形态不是冗余。Raw 服务审计与排障,A 是协议层的自然产物,Canonical 是四个下游消费者唯一该认的东西,Derived 是"没有原文时不许假装有原文"的产物。第三章的分配问题,本质上就是 A 和 Canonical 没有及时分开造成的。
三、67 次分配到底花在哪
先看三处已经复核过的事实。
第一处:入口重复拷贝。processMessage在插件路径上先做了一次make+copy,紧接着json.Marshal(pluginData.Data)又把整块覆盖掉。前半段是纯浪费。
// 修复前的形态(现已重排为先判插件、再决定是否拷贝)jsonData:=make([]byte,len(payload))copy(jsonData,payload)if!skipProtocolDecode&&preprocessor.NeedsProtocolPlugin(protocolName){pluginData,err:=plugins.ProtocolDecodeOnce(...)...jsonData,err=json.Marshal(pluginData.Data)// 上面那次拷贝白做了}第二处:设备日志重包一遍。直连路径原本把 A 交给BuildPropertyLogPayload重新构造成标准信封,再Encode成 JSON 存日志——一份内容、两次序列化。
第三处:影子写入时逐属性建节点。UpdateReportedAsync内部对每个属性调CreatePropertyNode,而 Canonical 本来是批量的,逐属性建节点的开销与属性数线性相关。
把这三处叠起来,再加上map[string]any从入口贯通到扇出点(每次取值都要类型断言、每层都可能再拷一次),就是malloc/s量级到百万那类现象的直接来源。根因不是某一处写错了,而是"同一个值在热路径上被当成不同东西处理了好几遍"。
四、扇出点只有一个
收敛的目标架构里,属性路径的扇出边界只有三个地方:
DeviceActor.flushBatch(正常异步路径)fallbackToSyncProcessing/ApplyPropertyEffects(Actor 投递失败后的同步降级)- 网关批量拆包之后进入上述边界(外层审计日志除外)
扇出之后,以下动作一律禁止:再跑一次ParseProperties、再对属性调CreatePropertyNode、把 A 喂给影子或告警、业务侧从日志 Content 里反解析 Canonical。
这条边界不是风格问题。扇出前做一次归一,得到的是"四个消费者语义一致";扇出后各转各的,得到的是"四个消费者各自看起来都对,但值不一样"。
五、口径分裂是什么样子
收敛之前的扇出是这样的:
| 消费者 | 吃什么 | 结果 |
|---|---|---|
| 影子 | A(原始 params) | 未过类型转换 |
| 属性告警(Direct) | A 解包出的毛值 | 未过ConvertValue |
| 属性告警(影子旁路) | 从影子 Reported 里再解一次 | 同一次上报被评第二遍 |
| TSD | Parsed != nil判断 | 无法区分"没解析"和"解析失败" |
| 审计日志 | Actor 路径重包 JSON / 同步降级存原文 | 同一设备存在两种 Content 形态 |
三件事值得单独说。
第一,属性告警是双路,而且没有去重。事件告警有tryMarkEventProcessed之类的时间戳去重,属性没有。同一条上报会走CheckPropertyAlarmDirect评一次,再经UpdateShadow → OnDeviceUpdate → checkPropertyAlarm评一遍。
第二,阈值判断吃的是毛值。影子路径只从 Reported 里提取{value, time},不做类型转换。这意味着温度属性在物模型里声明成浮点、设备上报字符串 “30.5” 时,影子看到的和告警比的是两个东西。
第三,日志 Content 两种形态并存。Actor 路径存的是重包后的合法 JSON,同步降级路径存的是string(data.PayLoad)原文。这不是未来的风险,是现网已经存在的状态——前端设备详情弹窗里JSON.parse没有 try/catch,遇到非 JSON 的 Content 会抛异常,弹窗打不开或显示旧值。
六、收敛一:把"解析状态"从约定变成字段
扇出要按"解析成没成"来决定给谁发数据,而原来的表达方式是靠约定:
Parsed = nil同时表示"还没解析"和"解析失败"Parsed = 空 map表示"解析成功但没有匹配到物模型属性"
一个字段承担三重含义,任何新写的消费者都得先读懂这个约定。所以第一步是把状态显式化:
// pkg/iotModel/device.gotypePropertyParseStatestruct{Attemptedbool// 是否已执行物模型解析Succeededbool// 解析过程是否成功(无 error)}三种组合的含义就清楚多了:
| Attempted | Succeeded | Canonical | 含义 |
|---|---|---|---|
| false | — | nil | 尚未解析,不应扇出业务 |
| true | true | 非 nil(可为空 map) | 成功;空 map = 没有匹配到物模型属性 |
| true | false | nil | 解析失败 |
有了状态,扇出策略就能写成一张确定的表:
| 状态 | 审计日志 | 影子 | 告警 | 北向 | TSD |
|---|---|---|---|---|---|
| 成功(有匹配) | 按 LogKind 写 | 更新 | 评估 | 发 | 写 |
| 成功(无匹配) | 按 LogKind 写 | 跳过 | 跳过 | 跳过 | 跳过 |
| 失败 | 写(便于排障) | 跳过 | 跳过 | 跳过 | 跳过 |
时序数据库读侧的兼容是这条改动里最容易踩的地方。改之前 TSD 遇到Parsed = nil会拿 Content 兜底再解析一次——当时能救回数据,是因为 Actor 路径的 Content 恰好是重包出来的合法信封 JSON。Content 一旦换成真实原文(可能是纯文本、也可能是二进制协议原文),这个兜底要么静默返回空、要么直接抛异常。所以顺序是死的:先有显式状态,再切日志形态;反过来做,切换期就会丢数据。
七、收敛二:属性告警从双路到单路
收敛后的合同写成一句话:每条属性上报,对告警只评估一次,入口只有CheckPropertyAlarmDirect,入参必须是 Canonical 加合并快照。
现在DeviceActor.flushBatch里是这样:
iflen(mergedCanonical)>0{mergedSnapshot:=alarmSrv.MergeCanonical(a.lastProps,mergedCanonical)a.lastProps=mergedSnapshotifhandler:=alarmSrv.GetAlarmShadowHandler();handler!=nil{handler.CheckPropertyAlarmDirect(alarmCtx,a.productKey,a.deviceKey,mergedCanonical,mergedSnapshot)}// ... 场景联动 ...updateShadowCanonicalFn(shadowCtx,a.productKey,deviceDetail,mergedCanonical)}注意最后的顺序:先评估告警,再更新影子。影子更新不再触发属性告警,OnDeviceUpdate被清空成 no-op:
// OnDeviceUpdate 设备影子更新通知。// Phase 3:属性/事件告警已收敛为 Actor Direct-only,此处不再二次评估。// 持续条件到点补评仍由 recheckSustainedConditionAlarms → checkPropertyAlarm 负责。func(h*AlarmShadowHandler)OnDeviceUpdate(deviceShadow*shadow.DeviceShadow)error{ifh.ctx.Err()!=nil{returnh.ctx.Err()}_=deviceShadowreturnnil}它没有变成"什么都不做"——持续条件的到点补评还在,只是改由兜底扫描recheckSustainedConditionAlarms负责,而不是靠每次影子更新顺带评一遍。这个区分很重要:影子更新触发评估是"数据驱动",兜底扫描是"时间驱动",后者才是持续条件(比如"连续 5 次超阈值"“5 小时未恢复”)真正需要的触发源——设备一直不上报时,数据驱动的评估根本不会被触发。
事件告警走同样的路:只留CheckEventAlarmDirect,时间戳去重从"必需"降级成"防止误双调的防御"。
八、收敛三:一个写了很久却没人读的字段
DeviceActor里有个lastProps,注释写着"跨批累积",实际状态是只写不读。跨批持续条件(“温度连续 5 次超过 30 度”)单靠当批 Canonical 根本没法评估,因为一次批量里只有一个时刻。
修法是把合并动作提到评估之前,并且明确归属:
// 1. 当批 Canonical 逐条解析出来,合并成 mergedCanonical// 2. 与历史合并:parsed 覆盖 lastProps 的同名键mergedSnapshot:=alarmSrv.MergeCanonical(a.lastProps,mergedCanonical)a.lastProps=mergedSnapshot// 3. 两个都交给唯一入口:当批新值用于单次阈值,合并快照用于持续条件与恢复判断handler.CheckPropertyAlarmDirect(ctx,productKey,deviceKey,mergedCanonical,mergedSnapshot)还有一个容易被忽略的细节:Actor 是按设备创建的,进程重启后内存里的lastProps就没了。所以 Actor 启动时要从影子里把历史灌回来:
func(a*DeviceActor)onStarted(ctx actor.Context){a.seedLastPropsFromShadow()}灌回来的数据必须走与 Canonical 一致的投影——影子存的是{time, value},读出来要组装成属性节点,不能在灌种子的路径上又引入一套"只取 value 不看类型"的老逻辑。否则只是把口径分裂从热路径搬到了启动路径。
九、收敛四:Raw 的所有权说不清楚
这一条是并发安全的账。
processMessage成功时返回的res是独立堆缓冲,可以安全跨异步持有;但错误分支返回的是入站原始[]byte。调用方目前不会误用错误分支的返回值,所以没有实际故障——可问题在于纯[]byte类型在编译期无法区分"owned"和"借用"。将来某个调用方拿着错误分支的返回值去异步用,编译器不会拦。
解法是让类型自己说话:
// Only err==nil 时返回非空 OwnedPayload(独立堆缓冲,可跨异步持有);// err!=nil 时返回空包装,调用方无法拿到入站借用缓冲。funcprocessMessage(...)(OwnedPayload,error){// ...returnOwnedPayload{},fmt.Errorf("plugin decode failed: ...")// ...returnOwnedPayload{data:jsonData},nil}同时,这条合同明确绑定了两件"不许做"的事:投递前不许对 PayLoad 做无谓的 Clone(既然已经是 owned,Clone 就是纯浪费),以及不许用unsafe做零拷贝字符串再跨异步持有。
十、收敛五:日志 Content 有三种语义
审计日志的 Content 到底是什么,原来靠"哪条路径写的"来推。现在写进枚举:
| LogKind | 场景 | Content 是什么 |
|---|---|---|
| Raw | 直连属性上报 | 入口 owned 的报文原文 |
| GatewayBatchAudit | 网关批量外层审计 | 网关报文的原文(只写一次) |
| DerivedEnvelope | 网关子设备业务日志 | 平台重编的 JSON,不是子设备原文 |
第三种是这张表里最关键的一行。网关子设备的属性是嵌在网关报文里的,子设备自己没有独立原文。原来这里会重包一份 JSON 当 Content 存下来,形式上和一手的原文没法区分——排障时看到一份"像原文的东西",其实是谁也说不清算谁的派生品。现在它老老实实标记成 DerivedEnvelope。
日志写入侧还顺手清掉了一处死分配:
ifstr,ok:=obj.(string);ok{// string 热路径直接用,禁止再 []byte(str) 死分配content=str}else{buf:=jsonBufferPool.Get().(*bytes.Buffer)// ...}日志内容切成 Raw 之后,string 分支就是热路径,而它的分配次数比原来"构造 map → Encode → 拷贝 → 转 string"那条路少一次。
十一、Raw 驻留不是免费的
把 Raw 一路带在内存里有个直接的代价:它活得比以前久。以前 PayLoad 在入口反序列化完就能释放,现在要挂在 Actor 的邮箱和批量缓冲上直到 flush。
这笔账可以估出来。固定条件下:批量大小 20,每设备入站上限 100,flush 间隔 100ms,单设备上报速率 10 msg/s,报文约 500 字节,活跃 Actor 1000 个:
(20 + 100 + 0.1 × 10) × 500B × 1000 ≈ 60.5 MB公式里的 100 是每个设备的入站闸门MailboxCapacity。它不是一个软限制——Send*占不到槽位就直接返回ErrActorMailboxFull,调用方必须走同步降级,不许假装投递成功:
varErrActorMailboxFull=errors.New("device actor mailbox full")这个闸门同时兼任内存保护:单设备的 in-flight 消息数被钉在 100 以内,Raw 驻留就有了上界。所以验收的时候,heap_inuse的增量要去和公式估的 60.5MB 对账,实测超过估算值 1.2 倍就算没达标。只看分配次数、不看驻留,等于把"少分配了但留久了"这种情况当成优化。
十二、验收为什么必须两个指标一起看
基线数据是这样:
| 项 | before | after |
|---|---|---|
| allocs/op | 67 | 37 |
| B/op | 5115 | 4811 |
| ns/op | ≈5140 | 3475–3923 |
分配次数降了约 45%,落在目标区间 −30% ~ −50% 之内。但列出来的验收项有五个:allocs/op、heap_inuse峰值、malloc 次数、GC 频次,以及日志分支的分配明细。
为什么不能只看一个。只报allocs/op,会漏掉"分配少了但对象留得久了"这类退化;只报heap_inuse,会漏掉高频小对象造成的分配密度问题。而且对比的两端必须是同一边界——基线端如果取的是池化 buffer 分支,目标端取的是非池化string(),两者不等价,算出来的改善就是假的。所以日志这一项要拆下来单独比:切换后的日志分支分配次数,必须不大于切换前 Encode 分支的次数,否则视为口径没对齐。
还有一条验收表述被明确禁止:不许把"四个下游最终状态一致"当成成功标准。影子队列可能熔断,TSD batcher 可能满而丢弃,告警异步提交可能失败,北向发布可能失败——四个 sink 各自独立。能承诺的只是"输入语义一致":四个消费者拿到或序列化出去的业务值,都来自同一份 Canonical。可靠性体现在另一处,即每个 sink 的入队成功、丢弃、超时、失败原因都要有计数。
告警 sink 原来是没有计数的。往细里数有四个失败点:提交池满、评估超时、评估异常、写库失败。池满这一处的处理方式是刻意的——立即返回错误而不是阻塞等待:
// 池满时立即返回错误而不是阻塞等待:调用方位于设备 Actor 消息循环内,// 阻塞会把背压传导到整条上报链路(邮箱积压 → inFlight 打满 → 降级同步处理)。func(h*AlarmShadowHandler)submitTask(taskfunc())error{select{caseh.workerPool<-struct{}{}:// ...default:returnfmt.Errorf("告警 worker pool 已满,丢弃本次评估,当前容量: %d",h.maxWorkers)}}丢了就丢了,由后续上报或持续条件兜底扫描补评。这是"宁可少评一次,不让背压打穿上报链路"的选择——评估失败的直接后果是漏一次告警,而阻塞的后果是整条链路降级。
十三、刻意留着没动的一处
收敛做完,链路里还剩一处明确没有覆盖的旧逻辑:pkg/dcache/deviceCacheData的最新值缓存,和internal/logic/analysis的分析页,仍然会对 Redis 日志列表里的 Content 做一次反解析。
现在它还能用,纯属巧合——直连路径的 Raw 恰好就是标准信封 JSON,能被解析器认出来。一旦 Raw 换成非信封形态(比如某类协议插件的原文),这两处会静默丢数据:解析返回空,被当成"已解析成功但没匹配",既不报错也不跳过。
这一处没有在这次收敛里解决,是因为它属于"读侧的历史包袱",改动面比写侧大。把它列出来而不是假装不存在,是因为它正好说明了这次收敛的边界:能保证的是四个已知消费者语义一致,不能保证的是所有下游都不再自己反解一遍。
十四、这次收敛不做什么
- 不改变 MQTT / Poller / Tunnel 三路径汇聚的模型,只是把汇聚之后的处理收紧。
- 不强制影子在 Redis 里的文档改成全新 schema,落盘仍然兼容
{time, value}。 - 不改前端日志列表去展示 Canonical——认定值看影子、时序库和曲线,日志列表保持它原本的用途。
- 不用
GOGC/GOMEMLIMIT来替代这套改造。调 GC 参数能让堆看起来平稳,但热路径上该转四遍还是四遍。
这条链路的改造没有引入新概念,做的是把已有的概念钉死:一个值在扇出前只归一一次;解析状态用字段表达而不是用 nil 表达;告警只认认定值;日志内容属于哪一类写进枚举;谁拥有那块内存由类型负责声明。剩下的都是这些约束自然带来的结果——包括从 67 降到 37 的那部分分配。
项目地址:
主仓库:https://github.com/sagoo-cloud
文档站:https://iotdoc.sagoo.cn