ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Metabase.mq:Metabase 的持久化消息队列——从事务性 Outbox 到 At-Least-Once 交付的完整设计

Metabase.mq:Metabase 的持久化消息队列——从事务性 Outbox 到 At-Least-Once 交付的完整设计 Metabase.mqMetabase 的持久化消息队列——从事务性 Outbox 到 At-Least-Once 交付的完整设计【免费下载链接】metabaseThe easy-to-use open source Business Intelligence and Embedded Analytics tool that lets everyone work with data :bar_chart:项目地址: https://gitcode.com/GitHub_Trending/me/metabasemetabase.mq是 Metabase 内置的持久化消息队列子系统为应用提供统一的发布/监听 API同时支持基于应用数据库Quartz 表 queue_message_outbox表的持久传输与纯内存传输两种后端。本文以 src/metabase/mq/README.md 为主体结合src/metabase/mq/下的源码实现完整讲解队列声明、四层发布语义、事务性 Outbox 的三阶段保证、并发限制策略、终端失败钩子与幂等性要求读完后可独立理解并在 Metabase 代码库中正确使用这套队列设施。定位与整体架构metabase.mq的核心承诺是无论队列使用哪个后端消息被发布当且仅当产生它的业务事务提交这一保证都成立——Outbox 表始终位于应用数据库app DB中后端只负责搬字节。官方文档给出的分层架构如下引自 READMEUser Code with-queue │ ▼ Publish Pipeline (mq.publish) collect in with-queue → deduplicate → route by :transactional mode: immediate → 100ms coalescing buffer (mq.publish-buffer) outbox → queue_message_outbox row, published after commit (mq.queue.outbox) │ ▼ Transport Dispatch (mq.transport) ← multimethod on channel namespace :queue → queue backends │ ▼ Backends quartz — one-shot Quartz job per batch (push; clustered JDBC JobStore) memory — LinkedBlockingQueue per channel │ ▼ Delivery Engine (mq.impl) worker pool → handle! → listener fn各层对应的源码模块架构层源码模块统一公共 APIdef-queue!/def-listener!/with-queue/putsrc/metabase/mq/core.clj发布管线缓冲、事务路由src/metabase/mq/publish.clj队列注册表与配置校验src/metabase/mq/queue/registry.clj事务性 Outboxsrc/metabase/mq/queue/outbox.clj合批缓冲时间窗src/metabase/mq/publish_buffer.clj传输分发src/metabase/mq/transport.clj队列后端协议src/metabase/mq/queue/backend.clj失败/丢弃策略含:on-errorsrc/metabase/mq/queue/impl.cljapp DB 查询Outbox 行、Quartz triggersrc/metabase/mq/db.clj恢复清扫 Quartz 任务src/metabase/mq/task/outbox.clj全局设置重试次数、后端选择src/metabase/mq/settings.clj公共 API 通过 core.clj 的potemkin/import-vars从mq.listener、mq.publish、mq.impl、mq.init、q.registry聚合导出core.clj 的 ns 文档 特别强调Single-consumer 指的是注册层面每个队列只有一个 handler而非消息只会被处理一次——交付仍是 at-least-once监听器必须幂等。快速上手声明、监听与发布最小用法分三步引自 README 的 Quick Start;; Declare the queue (its broker-side identity and properties). (mq/def-queue! :queue/my-task {:transactional :try}) ;; Register a listener (the consumer-side handler). The body receives a *vector* of messages. (mq/def-listener! :queue/my-task [messages] (doseq [msg messages] (do-work! msg))) ;; Publish a message (inside a transaction — delivers after commit) (t2/with-transaction [_] (mq/with-queue :queue/my-task [q] (mq/put q {:key value})))几个必须记住的约定队列名必须是:queue/*命名空间关键词。注册表用 malli schema 强制校验队列名必须是 namespace 为queue的关键词registry.clj L35-L38传输分发同样依据渠道关键词的 namespace 选择后端非:queue/*渠道会直接抛异常transport.clj L11-L26。所有队列配置都放在def-queue!上——:exclusive、:max-concurrent-batches、:max-batch-messages、:dedup-fn等因为它们在每个节点的发布时刻生效与监听器注册在哪里无关def-listener!只负责接线消费者侧 handler。监听器注册时若队列尚未声明会直接抛异常——拼写错误在启动时被捕获而不是第一次发布时。def-queue!与def-listener!分离的设计动机core.clj ns 文档队列独立于监听器存在发布者可以路由到代码库任何位置声明的队列监听器可以生活在另一台节点上。每个队列有且仅有一个声明点。def-queue!宏展开时调用claim-queue-declaration!记录声明所在 ns同一 ns 重载合法最新配置生效但另一个 ns 声明同名队列会在 load 时抛异常避免两个defmethod互相覆盖、谁生效取决于加载顺序registry.clj L160-L184。队列配置项全解def-queue!的 config 是一个closed map未知 key 直接报错由 malli schema:metabase.mq.queue/queue-config校验registry.clj L40-L54配置项类型含义:transactional必填:require/:try/:never控制发布与外层 DB 事务的关系见下文路由表:exclusiveboolean为 true 时全集群至多一个该队列的 batch 在途:max-concurrent-batches正整数或0 参 fn单节点同时交付的 batch 数上限软节流fn 每次检查时求值可绑定defsetting实时调整:max-batch-messages正整数batch 大小的软目标默认 100发布合批与消费切片共用:dedup-fnfnmessages - messages发布前对 batch 去重:on-errorfn终端失败钩子batch 耗尽重试被丢弃时调用源码层面的两个关键细节:exclusive与:max-concurrent-batches互斥同时声明会在注册时抛异常——schema 中的校验函数直接给出错误信息a queue cannot declare both :exclusive and :max-concurrent-batchesregistry.clj L50-L53。:max-batch-messages默认 100且是软目标不是硬上限异步合批缓冲在高负载下会乐于超过它因为发一个更大的在线 batch 比发几个小 batch 更划算所以监听器不能假设 batch 不超过该值registry.clj L15-L21。0 参 fn 形式的:max-concurrent-batches每次调用都会重新求值因此绑定 setting 的上限可以实时生效而不是在注册时冻结fn 抛异常会被捕获并退化为无上限返回 nil节流是尽力而为的宁可让工作流下去也不能卡死调度器registry.clj L99-L117。动态上限的官方示例README(mq/def-queue! :queue/exploration-query {:transactional :require :max-batch-messages 100 :max-concurrent-batches #(explorations.settings/explorations-worker-count)})发布语义四层递进的延迟投递with-queue被设计成读起来就像直接发布——写下消息就继续走。但在底层投递被刻意延迟了多层README 的 Publishing Semantics 一节共四条宏体必须成功。只有with-queue体正常返回消息才会入队异常会丢弃它们。对应实现run-with-bufferbody 成功才取buffer并路由异常则记录日志并重新抛出publish.clj L84-L104。外层事务必须提交。若调用处在t2/with-transaction块内直接或间接消息会一直持有到事务提交回滚即丢弃。这保证不会为一次未发生的数据库变更投递消息。合批时间窗。事务提交后消息进入时间窗缓冲同一渠道在滑动窗口内的连续发布会被合并成一个 batch 再落后端从而在突发流量下摊薄 per-batch 开销。若 flush 时 batch 够不到后端会被交给持久化 Outbox见下而不是在内存里反复重试直至丢弃——所以非事务性发布在app DB 可达的前提下能扛过后端故障。去重。派发前同 batch 内的重复消息被移除。源码印证publish_buffer.clj滑动窗口默认100ms*publish-buffer-ms*且每来一条新消息窗口就重置另有5000ms 的强制上限*publish-buffer-max-ms*自首条消息起超过即强制 flushL16-L23一个后台守护线程每 100ms 跑一次flush-publish-buffer!线程名mq-publish-buffer-flushflush 任务被try/catch Throwable包裹——因为scheduleAtFixedRate遇到未捕获异常会静默停止重排冻结整个 flusherL112-L136当某渠道缓冲量达到:max-batch-messages时会绊倒其 deadline让下一次 tick 排空它但发布者线程绝不阻塞在后端写入上L33-L63优雅停机时stop-publish-buffer-flush!会force?排空所有条目保证关闭不丢缓冲中的消息L138-L145。去重发生在发布入口publish!先取渠道的:dedup-fn应用之若有消息被丢弃会打点:metabase-mq/dedup-messages-droppedpublish.clj L21-L37事务性路径同样在 before-commit 插 Outbox 行前先做 dedupoutbox.clj L84-L93。立即发布的含义要精确理解:try在事务外 /:never的immediately只是马上交给发布管线仍要过 ~100ms 的合批缓冲上述第 3 步是在内存中短暂缓冲而非同步写后端。该窗口内崩溃会丢失非事务性发布——如果这不可接受就用:require或事务内的:try。事务性发布:transactional路由与 Outbox 三阶段每个队列必须声明:transactional模式它决定发布与外层 DB 事务的关系。路由逻辑实现于metabase.mq.publish/publish-collected!publish.clj L63-L82与 README 中的表格一一对应模式事务内事务外:require走 Outbox抛异常——事务是强制的:try走 Outbox立即发布*:never延迟到 after-commit但只在内存中持有无 Outbox 行立即发布*源码中:require事务外直接throw (ex-info Queue ... is :transactional :require and must be published inside a transaction.)publish.clj L78-L80。为什么需要 Outbox对:require/:try在事务内的发布业务写入提交后直接把消息交给后端是不安全的节点在业务事务提交之后、发布之前崩溃消息就丢了。事务性 Outbox模块metabase.mq.queue.outbox表queue_message_outbox位于 app DB补上这个缺口分三个阶段outbox.cljbefore-commit——为每个渠道收集的消息先 dedup、按:max-batch-messages切块、编码作为queue_message_outbox行在仍然打开的业务事务内部插入因此与业务写入原子提交before-commit 回调若抛异常整个事务回滚。实现为insert-outbox-rows!把插入行记回事务状态的::rowsoutbox.clj L77-L94。after-commit——每条插入的行发布到后端成功后删除直通无缓冲。每行独立 try/catch发布失败后端宕机、调度器不可用……该行原地保留等待恢复清扫只有发布成功的行才会被删除publish-outbox-rows!outbox.clj L96-L115。recovery——周期性清扫recover-outbox!重发任何超过 ~1 分钟recovery-age-ms 60s的遗留行由 src/metabase/mq/task/outbox.clj 中的 Quartz 任务驱动cron 每 5 分钟DisallowConcurrentExecution。它用FOR UPDATE SKIP LOCKED认领行Postgres/MySQLH2 用普通FOR UPDATE使不同节点上的并发清扫拿到不相交的行outbox.clj L63-L69。发布失败几乎总是瞬时的后端/DB 连通性所以失败行永不丢弃publish_attempts递增next_attempt_at按指数退避调度——1 分钟起步、翻倍、封顶 10 分钟recovery-retry-base-ms60srecovery-retry-max-ms600s退避公式见 retry-delay-ms L56-L61一直重试到发布成功。几个值得注意的源码细节清扫按id 键集分页每事务一页页大小recovery-page-size 50outbox.clj L41-L43清扫中若识别出后端整体不可用异常backend-unavailable?异常标记定义于 backend.clj L57-L69会立即停止本轮清扫reduced留待下次调度而不是把剩余每一行都 bump 一遍去锤一个已知不可用的后端outbox.clj L162-L180查询到期行的条件next_attempt_at为空且创建早于now - recovery-age-ms或next_attempt_at nowdb.clj L32-L50一个后来从代码中移除的队列的遗留 Outbox 行不是特例发布到后端不检查队列是否存在行照常发布并删除产生的 trigger 则像其他无监听消息一样由队列 reapermetabase.mq.task.queue-reaper随时间淘汰。净保证消息被发布当且仅当产生它的业务事务提交——无论队列使用哪个后端Outbox 表始终在 app DB。交付仍是at-least-onceafter-commit 窗口内崩溃会让清扫重发因此监听器必须幂等。:never跳过 Outbox 的快路径但不承诺可丢:never在快路径上跳过 Outbox事务内仍等待提交回滚即丢弃但只在内存持有——无 Outbox 行高负载下保持低成本就这一点它是best-effort——flush 前崩溃会丢消息。但它对后端可用性不是best-effort若 flush 够不到后端:never的 batch 会像:try一样回落到 Outbox后端故障不丢消息。选择:never的理由应该是想省掉每次发布的 Outbox 写而不是愿意丢消息。发布时 Outbox 回落非事务性发布未走 Outbox 路由的发布——事务外的:try、任何:never——经过合批缓冲后直连后端。若后端写失败调度器宕机、瞬时 DB 抖动缓冲不自管重试/丢弃它通过outbox/insert-batch!把整个 batch 写入queue_message_outboxnext_attempt_at留空如同崩溃遗留行由同一个恢复清扫接管——该行要老化超过 ~1 分钟recovery-age-ms才可能被捡起之后按退避重试、永不丢弃。因此交接后的 batch 延迟至多恢复窗口 清扫间隔不会立即发布但也不会丢。实现见handoff-to-outbox!publish_buffer.clj L65-L76。唯一的剩余丢失场景是 Outbox 插入本身也失败——即 app DB 也宕机了全面故障此时打日志并打点batches-dropped{reasonoutbox-handoff-failed}。另外注册表文档给出了性能注记registry.clj L245-L250事务性发送不走滑动时间窗缓冲消息在事务边界立即成行无法跨事务合并若消息有风暴风险可考虑:never换性能。消息序列化JSON 单向编码消息必须是 JSON 可序列化的。整个链路中batch 在发布边界编码一次mq.transport/publish!→mq.payload/encodetransport.clj L28-L34在交付边界解码一次mq.impl/deliver!→mq.payload/decode中间环节后端只搬动不透明字符串、从不窥视内容——所以编码不是后端的职责每个后端交付的消息形态完全一致。解码时 map 的 key 会被 keyword 化但值按 JSON 类型系统直通keyword 值变字符串、set 变 vector、日期变字符串。发布前还有一道check-serializable!校验不可 JSON 往返的值如异常对象会在调用点直接抛错而不是在线路上被静默损坏后永远无法投递payload.clj L35-L63在 publish.clj L95-L99 中于任何路由之前执行。队列语义保证交付不保证顺序Queue交付每个 listener at-least-once顺序无。见下重试有——最多queue-max-retries次默认 5后备存储QuartzQRTZ_*表事务性发布用queue_message_outbox积压持续存在直到被处理或耗尽适用场景工作必须可靠完成queue-max-retries是内部 setting默认 5settings.clj L6-L12queue-backend默认quartz当任务调度器被禁用MB_DISABLE_SCHEDULER时回落到memorysettings.clj L14-L21。保证交付是保证顺序不是。batch 并发交付、带退避重试失败的 batch 会落到更晚发布的 batch 之后且 Quartz 即使串行化 trigger 也不按提交顺序触发。不要写假设消息按发布顺序到达的监听器——也不要从通过的测试反推该保证因为内存后端可能碰巧看起来是 FIFO。若队列需要串行交付用:max-concurrent-batches 1每节点或:exclusive true全集群明确表达——并记住即便它们也只给你互斥不给你顺序。限制并发:exclusive与:max-concurrent-batches默认情况下队列的 batch 并发运行只受 Quartz 线程池约束——而该线程池与 sync、pulses 等所有定时任务共享。一个能 fan-out 的队列每行一条消息、每张图表一条……会毫不客气地吃光这个池。两个旋钮配置界限硬性执行者:exclusive true全集群1 个 batch是后端各自实现:max-concurrent-batches n每节点~n 个 batch否——节流trigger 获取 /fetch!空闲槽位计数都不设无界直到 Quartz 线程池——要点README 的 Limiting concurrency 一节 源码印证:exclusive是后端保证每个后端各实现各的Quartz 用 job 类上的DisallowConcurrentExecution内存后端则是队列已有在途 batch 时绝不再取第二个。共享的轮询驱动器不知道这个标志新后端若不实现就会静默失效。二选一不可叠加。同时声明在注册时抛异常:exclusive已是最严限制其上再叠 per-node 上限永远不可能收紧且两者由完全不同的机制执行。想要每节点大致同时一个用:max-concurrent-batches 1要精确同时一个用:exclusive。到达上限的节点是停止接收而非内部排队它报告自己当前处理不了该队列亲和性委托affinity delegate把这些 trigger 留在共享存储中WAITING等有容量的节点来取——这与节点不为没有监听的队列取消息是同一套机制。这也是真实 broker 唯一能兑现的限制形态现在最多取 N 个对应 SQS 的MaxNumberOfMessages、Redis 的 multi-pop 计数所以轮询式后端通过向fetch!询问 per-queue 空闲槽位来表达queue-free-slotsmap 中计数为nil表示未设上限无界已满的队列根本不在 map 里计数永不为零backend.clj L18-L28。:max-concurrent-batches是节流不是保证。两个执行点都与它把关的工作存在竞争Quartz 在 job body 运行前就由调度线程决定取什么轮询后端可能在槽位刚被占满的一刹那交出一个刚取回的 batch——所以单节点可能超出上限一两个 batch。此时 batch照常被交付不回退一条重新入队路径需要自己的退避、自己的指标还要防范上游节流失效时的热循环全为削一个小、有界、能自排的超额不划算。上限必须真的守住时用:exclusive。队列配置在所有后端上含义一致——这是契约而非巧合后端只允许在机制上不同Quartz 在 trigger 获取时推并过滤轮询后端向fetch!要空闲槽位计数绝不在配置含义上不同。新增后端必须达到这条标准。终端失败:on-error没有死信队列。batch 耗尽queue-max-retries后被丢弃没有 handler 时唯一的痕迹是一行日志和一个batches-dropped{reasondelivery-exhausted}计数器。对 fire-and-forget 的工作没问题但会无声地搁浅生产者留下的东西——比如一个停在pending状态等 handler 的行会永远停在那里。声明:on-errorhandler 来持久化记录终端失败README 的示例(mq/def-queue! :queue/run-report {:transactional :require :on-error (fn [{:keys [channel messages error attempts]}] (doseq [{:keys [report-id]} messages] ;; Only fail a report that is still pending. See at-least-once below: this ;; handler can fire for a batch a peer already completed. (t2/update! :model/Report {:id report-id :status pending} {:status error :error_message (ex-message error)})))})语义约束源码印证于 queue/impl.clj L13-L59只在终端丢弃时调用——不是每次失败尝试都调也不会为最终重试成功的 batch 调用handler 抛出的异常被记录并吞掉batch 反正要被丢弃坏 handler 不能把队列卡死在反复重投已耗尽的 batch上run-on-error!的 try/catch:on-error和交付本身一样是 at-least-once写成幂等的且容忍为实际已成功的工作触发——所以上例有:status pending守卫。为什么会为已成功的 payload 触发因为消息没有身份重试是一个带新 id 的全新 trigger尝试计数藏在其内部所以同一 payload 的两次交付是两个独立的 batch、两份独立的重试预算、没有任何东西把它们关联起来。若重复交付的那份失败而原件已成功失败副本会耗尽自己的预算、为早已完成的 batch 触发:on-error。这不是留给崩溃的角落案例Quartz 集群故障转移会在首个节点错过心跳时把 batch 重新触发到第二个节点上——而存活的节点在长 GC 暂停或 app-DB 延迟下就会错过心跳它永远不会被告知停止它的副本会继续跑。测试metabase.mq.queue.on-error-test钉住了这一行为duplicate-delivery-fires-on-error-for-an-already-succeeded-payload-test对应 test/metabase/mq/queue/on_error_test.clj要彻底关闭它需要消息身份 已交付记录。该钩子位于共享的重试 vs 丢弃策略metabase.mq.queue.impl/handle-batch-failure-policy!而不是任何一个后端里因此每个后端都继承它——轮询后端经共享轮询驱动器获得推式后端Quartz直接调用该策略。后端无法绕过这个策略实现重试预算所以未来新增的任何后端都免费带上:on-error。幂等性监听器的两条军规队列监听器必须幂等且必须与另一个节点上的自己并发运行也安全。要记住两件事README 的 Idempotency 一节batch 可能被重投。部分失败后它回到pending同样的消息会再来一遍batch 可能同时被交付两次。Quartz 集群故障转移在首个节点错过心跳时把 batch 重触发到第二节点——存活节点在长 GC、app-DB 延迟或连接池耗尽下会错过心跳且它永远不会被打断所以两份副本都会跑。另外未声明:exclusive的队列根本没有跨节点互斥N 节点 × Quartz 线程池可以同时处在同一个 listener 里。因此我做过这件事了吗的检查单独不够——两份副本可以同时检查、同时看到没有、同时继续推进。在乎的地方要让写入安全而不是让检查安全把更新条件化于触发消息的状态如{:id x :status pending}、用唯一约束做预约、或在事务内取行锁后再检查。测试工具with-test-mq的{:duplicate-delivery? true}选项会把每条消息交付两次用于揪出假设 exactly-once 的监听器。测试with-test-mq与内存后端用metabase.mq.test-util的with-test-mq针对隔离的内存后端跑测试队列声明放在命名空间顶层这样会注册进每次运行的全新注册表然后在绑定向量后紧跟 listener map 传入opts 里加{:duplicate-delivery? true}即可把每条消息交付两次来证明监听器幂等test/metabase/mq/test_util.clj(mq/def-queue! :queue/my-task {:transactional :try}) (deftest my-test (let [processed (atom [])] (mq.tu/with-test-mq [ctx] {:queue/my-task (fn [messages] (swap! processed into messages))} (mq/with-queue :queue/my-task [q] (mq/put q {:id 1})) (mq.tu/eventually! ctx #( 1 (count processed))) (is ( [{:id 1}] processed)))))相关测试入口test/metabase/mq/ 下的 backend_parity_test.clj后端行为一致性、publish_buffer_test.clj合批缓冲、publish_test.clj、listener_test.clj、payload_test.clj、quartz_affinity_test.clj 及 queue/on_error_test.clj。关键保证速查问题答案消息何时被发布:require/事务内:try当且仅当业务事务提交Outbox 原子写入事务外:try/:never~100ms 合批窗口内崩溃可能丢失后端故障会丢消息吗不会——batch 交接给queue_message_outbox恢复清扫接管仅当 app DB 同时不可用时才丢弃并打点outbox-handoff-failed恢复清扫多快行老化 60s 才可见Quartz 任务每 5 分钟一轮失败退避 1m 翻倍至 10m永不丢弃交付语义at-least-once无顺序保证监听器与:on-error都必须幂等并发上限:exclusive集群级硬性互斥与:max-concurrent-batches节点级软节流二选一超限 batch 照常被交付而非回退重试预算queue-max-retries默认 5耗尽即丢弃并触发:on-error如有【免费下载链接】metabaseThe easy-to-use open source Business Intelligence and Embedded Analytics tool that lets everyone work with data :bar_chart:项目地址: https://gitcode.com/GitHub_Trending/me/metabase创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表