尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

ThingsBoard消息优先级会怎么“插队“:一次源码走读

ThingsBoard消息优先级会怎么“插队“:一次源码走读 ThingsBoard消息优先级会怎么插队一次源码走读【免费下载链接】thingsboardOpen-source IoT Platform - Device management, data collection, processing and visualization.项目地址: https://gitcode.com/GitHub_Trending/th/thingsboard凌晨三点设备离线告警在页面上迟了几秒才出现同一时刻历史遥测数据在稳定写入。翻日志队列没有丢消息消费者线程也没卡死。为什么关键消息会在噪音里变慢答案藏在 ThingsBoard 消息优先级的实现里——每个 Actor 邮箱守着两条长度不受限的队列而插队权就写在一个不到二十行的轮询方法里。消息优先级到底解决什么痛点关键消息为什么不能排队一个 Actor 对应一个租户或一条规则链的处理上下文内部单线程串行消费天然无法并行。如果一条控制指令或规则链变更消息排在几万条遥测后面用户操作的可感知延迟就完全不可控。ThingsBoard 没有引入复杂的调度算法而是把急件和平件物理分开存让轮询逻辑永远先看急件箱。ThingsBoard 消息优先级落在哪从消费行为倒推两条队列先说最终生效的表现只要高优队列里有一条消息正常队列就被整个跳过。把这个现象倒着追回去。最底层是 TbActorMailbox.java 里的两条ConcurrentLinkedQueueprocessMailbox()每轮按actorThroughput配置项控制单次批量处理的消息条数拉取永远先 poll 高优队列// TbActorMailbox两条队列 一个轮询循环 private final ConcurrentLinkedQueueTbActorMsg highPriorityMsgs new ConcurrentLinkedQueue(); // 急件箱 private final ConcurrentLinkedQueueTbActorMsg normalPriorityMsgs new ConcurrentLinkedQueue(); // 平件箱 private void processMailbox() { boolean noMoreElements false; for (int i 0; i settings.getActorThroughput(); i) { // 一批最多处理 N 条 TbActorMsg msg highPriorityMsgs.poll(); // 先取高优 if (msg null) { msg normalPriorityMsgs.poll(); // 没有才轮到普通 } if (msg ! null) { actor.process(msg); // 交给 Actor 本体处理 } else { noMoreElements true; break; } } // 没取完则立刻再跑一轮 processMailbox取完才释放 busy 状态 }入口在tellWithHighPriority()调用方通过 DefaultTbActorSystem.java 按 ActorId 找到邮箱把消息投进急件箱并触发一次消费尝试。往上追调用面集中在 AppActor.java 这类中枢 Actor 里——分区变更、规则节点更新、算子字段状态恢复都是走高优通道的典型场景。而 common/queue/ 模块解决的是节点之间的消息传递Kafka Topic、Producer/Consumer 抽象与这里的邮箱优先级是两套独立的机制读源码时容易混淆。一条高优消息的完整链路长什么样高优和平优共用同一个轮询循环差别只在于 poll 的先后顺序而不是各自独立的线程。也就是说插队能力来自先看哪个箱子不来自额外的执行资源。这意味着高优消息能插到队头但插不了处理线程。高优通道背后的两个坑现象到根因一次说清⚠️ 坑一普通消息饥饿。 现象某租户遥测长时间不消费监控看到消费者线程一直是忙的但普通队列长度只增不减。 根因processMailbox()没取完一批就立刻重新投递自己如果高优消息的到达速率持续高于单批处理速率循环永远停在highPriorityMsgs.poll()那一步普通队列一次都轮不到。 规避审计tellWithHighPriority的调用面批量数据一律走tell()线上用两条队列的积压长度做对比告警高优积压长期大于零就要介入。⚠️ 坑二销毁中的 Actor 被复活拖进循环。 现象某个初始化失败的规则节点 Actor 反复打印重试日志节点更新操作却毫无效果。 根因enqueue()里有个特判——投递中如果 Actor 已销毁普通消息直接走onTbActorStopped收尾但RULE_NODE_UPDATED_MSG这类高优消息会把destroyInProgress复位、重新initActor()。这本来是节点被改了配置老 Actor 作废重建的修复手段但如果初始化失败的根因没消除每次更新都会触发一轮销毁—复活—再失败。 规避把复活机制当作兜底而不是修复出现 INIT_FAILED 日志时先解决初始化失败本身依赖数据缺失、配置错误而不是频繁发节点更新去踢它。接下来你可以做的事 给两条邮箱队列加监控在processMailbox出口读highPriorityMsgs.size()和normalPriorityMsgs.size()接进 Prometheus 后按 ActorId 出图饥饿问题在曲线上就是一眼的事。 把tellWithHighPriority的调用点列成清单在仓库里全局搜这个方法名逐个确认这条消息真的急吗多数遥测类路径误用高优通道都能在这里暴露出来。 自定义规则节点时先定通道策略节点自身发出的下游消息走tell还是tellWithHighPriority写进节点设计文档避免团队各写各的。顺带一提ThingsBoard 消息优先级 的监控指标怎么设计是这套机制落地时最常被问到的问题。【免费下载链接】thingsboardOpen-source IoT Platform - Device management, data collection, processing and visualization.项目地址: https://gitcode.com/GitHub_Trending/th/thingsboard创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表