ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

EMQX 客户端消息分页 API 修复:max_payload_bytes 截断后 meta.position 游标语义修正

EMQX 客户端消息分页 API 修复:max_payload_bytes 截断后 meta.position 游标语义修正 后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载导读本篇文章围绕 EMQX 管理 API 中GET /clients/{clientid}/mqueue_messages与GET /clients/{clientid}/inflight_messages两个接口的分页游标meta.position语义展开当单页响应因max_payload_bytes默认 1MB被截短时旧实现返回的position会指向被截掉的消息之后导致下一页跳过这些消息、造成消息看起来丢失的假象本次修复将position回退到页内最后一条实际返回的消息使下一页严格从被截断处继续。读完本文你将掌握这两个接口的完整参数语义、游标分页的底层实现消息队列/飞行中窗口的排序键与截断算法以及如何在实战中正确翻页而不重不漏。一、两个接口的业务背景GET /clients/{clientid}/mqueue_messages与GET /clients/{clientid}/inflight_messages是 EMQX 客户端管理 API 的一部分分别用于查询某个客户端会话的两类消息mqueue消息队列离线期间或消费速度跟不上时积压的 QoS 消息按mqueue_priority与入队时间排序底层由 emqx_mqueue.erl 实现inflight飞行中窗口已发送但尚未收到 PUBACK/PUBREC 确认的 QoS 1/2 消息按入队时间排序底层由 emqx_session_mem.erl 的inflight_query/2实现。在 HTTP 层两个接口由 emqx_mgmt_api_clients.erl 注册路由其处理函数分别落在mqueue_msgs/2与inflight_msgs/2源码位置统一走list_client_msgs/3的分页管线。二者的 OpenAPI schema 由client_msgs_schema/4生成源码位置响应结构一致仅mqueue_messages的消息条目比 inflight 多出mqueue_priority字段。二、请求参数与响应结构源自源码 schema两个接口的查询参数在client_msgs_params/0中统一定义源码位置参数类型默认值说明clientidstring必填path客户端 ID位于 URL 路径中payloadnone/base64/plainbase64消息 payload 的编码方式max_payload_bytesbytesize1MB单页所有消息 payload 的总字节上限必须 0否则返回参数校验错误positionstring无从头开始游标位置由上一页响应的meta.position返回用于翻页limitinteger100单页最多返回的消息条数默认值取自 emqx_mgmt.hrl 的DEFAULT_ROW_LIMIT并经 emqx_mgmt_api.erl 兜底max_payload_bytes的正数校验由max_bytes_validator/1保证源码位置测试用例也对max_payload_bytes0、-1、-1MB、0MB等非法输入做了 400 响应断言测试位置。成功响应的 JSON 结构为{ meta: { start: 第一页第一消息的位置后续每页恒定, position: 游标下一页应携带的 position读尽时为 end_of_data }, data: [ { msgid: ..., qos: 1, topic: ..., publish_at: 1234567890, from_clientid: ..., from_username: ..., inserted_at: ... } ] }其中mqueue_messages的每条消息还额外带有mqueue_priority字段字段定义见 源码位置字段断言见 测试位置。当payloadbase64时payload字段为 Base64 编码字符串payloadplain时直接返回原始 payloadpayloadnone时不返回 payload 字段此时单条消息格式化字节数为 0见 format_payload/3。三、游标分页的底层实现limit 截断 vs payload 截断理解本次修复必须先厘清这条管线的两级截断。3.1 第一级按 limit 截断游标生成list_client_msgs/3源码位置向会话进程请求limit 1条lookahead若实际取回多于limit条说明后面还有数据此时take_client_msgs_page/4把meta.position设为第 limit 条消息的位置若不足limit条则设为end_of_data源码位置。游标的具体形态取决于消息类型msg_position/2mqueue位置是二元组{mqueue_insert_ts, mqueue_priority}编码为字符串insert_ts_priority优先级为无穷大时写作infinity解码见 decode_mqueue_pos/1inflight位置是单个整数inflight_insert_ts。会话侧据此定位续页起点mqueue 队列在 emqx_mqueue.erl 的query/2中通过skip_until/2跳过insert_ts(Msg) MsgPos的消息源码位置即下一页从严格大于游标位置的消息开始inflight 窗口在 emqx_session_mem.erl 的inflight_query/2中用lists:dropwhile(fun(M) - inflight_insert_ts(M) Position end, ...)同样丢弃游标处及更早的消息。这意味着只要 position 指向某条真实消息下一页就不会漏、也不会重——这正是修复能依赖的语义基础。3.2 第二级按 max_payload_bytes 截断格式化阶段拿到会话侧返回的消息列表后format_msgs_resp/4进入格式化阶段源码位置。format_msgs/4源码位置用foldl_while逐条计算格式化后的字节数并累加第一条消息无条件返回——即使它单独就超过max_payload_bytes也会包含在响应中测试对此有专门覆盖见 测试位置从第二条起一旦SizeAcc PayloadSize MaxBytes立即停止累加并把此时本应返回但被截掉的最后一条候选消息标记为{truncated, LastMsg}若总字节数未超限则标记为all。3.3 Bug 根因两级截断的游标不一致修复前的问题在于meta.position由第一级 limit 截断生成指向 limit 条消息的最后一条而页面最终展示的内容却可能被第二级 payload 截断提前斩断。当max_payload_bytes把页面截短时返回的position指向的是被 payload 截掉的那批消息之后的位置即仍在消息序列中更靠后的地方于是客户端拿着这个position请求下一页时会话侧会跳过被截掉的消息——这些消息从未被返回却又被翻页跳过了表现得就像消息丢失。典型症状客户端详情的mqueue_len计数明显大于两个接口累计返回的消息数mqueue_len等字段定义见 源码位置逐页翻完后仍有多条消息不可见。四、修复方案position 回退到页内最后一条返回消息本次变更fix-18509.en.md的修复核心在adjust_pager_position/3源码位置若格式化阶段未发生 payload 截断allposition保持第一级 limit 截断生成的值若发生了 payload 截断{truncated, LastMsg}则将meta.position改写为页内最后一条实际返回消息的msg_position(MsgType, LastMsg)。结合 3.1 节会话侧的严格大于游标语义position指向最后一条已返回消息后下一页恰好从被截掉的第一条消息继续从而实现逐页翻页不重不漏。这条逻辑与源码注释完全一致源码注释当max_payload_bytes截短页面时将续页位置回退到最后一条返回消息使下一页从被截断的首条消息开始。注意该修正发生在格式化之后因此同时适用于 mqueue 与 inflight 两个接口由于 position 是基于真实消息mqueue_insert_ts/inflight_insert_ts而非基于返回条数的偏移量无论 payload 编码是base64、plain还是none修正逻辑都保持一致。五、测试用例的完整验证修复行为在 emqx_mgmt_api_clients_SUITE.erl 中有非常完整的回归覆盖t_mqueue_messages/1与t_inflight_messages/1共用test_messages/6的骨架测试位置构建数据client_with_mqueue/3创建 MQTT v5 客户端订阅 QoS1 主题后断开消息转入 mqueueclient_with_inflight/3用auto_acknever让消息停留在 inflight 窗口测试位置全量断言不带max_payload_bytes限制时单页返回全部消息position为end_of_data且start指向第一条消息截断场景本修复的核心断言payloadencmax_payload_bytes1limit2时页面仅含第一条消息并断言TruncatedPos等于第一条消息的位置而不是被截掉消息的位置同时TruncatedPos不等于end_of_data测试位置——这正是position 回退到最后一条返回消息的直接验证逐页接力翻页从TruncatedPos起以limit2max_payload_bytes1逐条翻页断言每一页恰好返回 1 条、payload 依次为2,3,...,Count直到最后一页position变为end_of_data测试位置。若修复缺失第 2 页起就会开始跳消息此断言必然失败非法参数positionnot-int、limit-5等返回 400测试位置不存在客户端返回 404测试位置。另外从 测试代码 可以看出t_inflight_messages同时被列入持久会话persistent session测试组即该接口对持久会话场景同样要求分页正确性。六、实战如何正确翻页消费全部消息基于上述语义推荐的分页模式是标准的游标接力首次请求不携带position从第一页开始每页响应后取出meta.position若为end_of_data翻页结束否则将该值原样作为下一次请求的position参数重复直到end_of_data。以 curl 为例假设管理 API 监听在 18083 端口、Dashboard 账号为admin:public# 第一页从会话消息队列头部开始 curl -u admin:public \ http://127.0.0.1:18083/api/v5/clients/my-client/mqueue_messages?payloadbase64limit50 # 后续页把上一页 meta.position 的值填入 position curl -u admin:public \ http://127.0.0.1:18083/api/v5/clients/my-client/mqueue_messages?payloadbase64limit50position上页positioninflight 接口的用法完全相同仅替换路径中的mqueue_messages为inflight_messages。需要强调的实战要点务必原样使用响应中的position不要自行推导或偏移它是基于消息时间戳与优先级生成的不透明游标mqueue 为insert_ts_priority形式inflight 为整数形式max_payload_bytes只约束 payload 总量不保证返回条数极端情况下页面可能只返回 1 条首条 payload 超限时仍返回这是源码保证的行为不要据此判断数据是否取完应以position end_of_data为准响应中的meta.start全程恒定指向队列头部第一条消息可用于核对首条消息的一致性但不应用于翻页若追求完全一致的响应体大小可将payloadnone与max_payload_bytes组合使用若需保留 payloadbase64是默认且推荐的编码plain对二进制 payload 可能因 JSON 序列化失败而返回INVALID_PARAMETER见 format_msgs_resp/4。七、小结本次修复fix-18509.en.md解决的是 EMQX 客户端消息查询接口中一个隐蔽且易被误判为丢消息的分页缺陷max_payload_bytes把单页截短后旧实现返回的meta.position越过被截消息导致续页跳读。修复通过 adjust_pager_position/3 将游标回退至页内最后一条已返回消息并借助会话侧严格大于游标的取数语义emqx_mqueue.erl 的skip_until/2与 emqx_session_mem.erl 的sublist_from_pos/3保证了逐页翻页的不重不漏。测试套件 emqx_mgmt_api_clients_SUITE.erl 用max_payload_bytes1的逐条翻页断言固化了这一行为为依赖这两个接口做消息审计、离线消息补拉或飞行中消息巡检的集成方提供了可靠契约。赞分享后端物联网消息队列通信【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址https://gitcode.com/gh_mirrors/em/emqx点击查看免费下载相关推荐EMQX 端到端追踪修复解析WebSocket 客户端出站消息为何失联及修复原理EMQX 端到端追踪修复解析WebSocket 客户端出站消息为何失联及修复原理 本篇文章基于 EMQX 仓库变更记录 changes/ee/fix 14后端物联网消息队列通信Eclipse Mosquitto 1.6.12 发布详解QoS 2 消息内存泄漏修复与客户端退出码修正Eclipse Mosquitto 1.6.12 发布详解QoS 2 消息内存泄漏修复与客户端退出码修正 导读 本文围绕 Eclipse Mosquitto后端消息队列消息路由EMQX 修复 Prometheus emqx_messages_retained 指标恒为 0让保留消息Retained Message写入量真正可观测EMQX 修复 Prometheus emqx_messages_retained 指标恒为 0让保留消息Retained Message写入量真正可观测后端物联网消息队列通信上一篇Apache Maka 的 ACP v1 实时会话行为订阅保留、流式投影与取消关闭语义详解下一篇如何在OpenMind框架中使用TinyMistral-248M-openmind完整代码示例创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表