Skip to main content

Crate dtmrs

Crate dtmrs 

Source
Expand description

§dtmrs

Rust 写的分布式事务管理器,对标 DTM。 支持 SAGA / TCC / 二阶段消息 / XA / workflow 五种模式, 存储可跑 sqlite / postgres / mysql,对外有 HTTP 和 gRPC 两套等价接口。

这个 crate 是门面:把各层重新导出到一处,省得挨个加依赖。 想精确控制依赖的话也可以直接用底下的 dtmrs-core / dtmrs-server / dtmrs-barrier / dtmrs-xa

§两种用法,别装错

dtmrs 的两端要的东西完全不同,feature 就是按这个切的:

你的角色装什么拿到什么
跑协调器(TC)dtmrs(默认)Embeddeddtmrs 二进制
业务服务(RM)dtmrs = { features = ["barrier"], default-features = false }barrier —— 幂等所必需
业务服务 + XA再加 "xa"xa

业务侧只需要屏障,不该把整个协调器和 axum/tonic 拖进去 —— 所以 barrier 不在默认 feature 里,default-features = false 之后它非常轻。

§嵌入式:不用单独部署服务

这是相对 DTM 的结构性差异 —— 协调器当库链进你自己的进程, 分支可以直接是进程内的函数:

use dtmrs::{BranchResult, Embedded};

let tc = Embedded::builder("sqlite:app.db")
    .handler("扣款",     |_ctx| async { BranchResult::Success })
    .handler("扣款撤销", |_ctx| async { BranchResult::Success })
    .start()
    .await?;

tc.saga("order-1001")
    .step("local://扣款", "local://扣款撤销")
    .step("http://shipment/create", "http://shipment/cancel")  // 可跟远端混用
    .submit()
    .await?;

也可以当独立服务跑:cargo install dtmrs 之后 DTMRS_DB=sqlite:dtmrs.db dtmrs

§业务侧接入:屏障不是可选项

分支接口一定会被重复调用(TC 重试 + 崩溃恢复),所以必须幂等。 把屏障记录和业务 SQL 放进同一个本地事务

use dtmrs::barrier::{BranchBarrier, Decision};
use dtmrs::Backend;

let mut bb = BranchBarrier::new(Backend::from_url(&url), tt, gid, branch_id, op)?;
let mut tx = pool.begin().await?;
if bb.decide(&mut tx).await? == Decision::Execute {
    // 业务 SQL —— 必须在这个 tx 里
}
tx.commit().await?;   // 原子性的来源

§一条写错就会数据不一致的规矩

超时 ≠ 失败。 5xx、连接超时、gRPC 的 UNAVAILABLE / DEADLINE_EXCEEDED 都表示结果未知 —— 对方可能已经成功了,这时候回滚就是不一致。 只有 HTTP 409、gRPC ABORTED、或响应体里的 FAILURE 才触发补偿。 见 BranchResult

Modules§

barrier
子事务屏障 —— 业务侧(RM)做幂等用,接入 dtmrs 的必需品
dialect
方言层 —— 一套 SQL 模板跑 sqlite / postgres / mysql。
server
协调器(TC)本体:推进器、嵌入式门面、HTTP/gRPC 接口
store
存储层:一套 SQL 跑 sqlite / postgres / mysql
xa
XA 两阶段提交的业务侧(RM)助手

Structs§

BranchBarrier
BranchCtx
分支被调用时拿到的上下文。业务侧用它做幂等(配合 dtmrs-barrier)。
Embedded
EmbeddedBuilder
RetryPolicy
重试退避策略。
SagaStep
一个 SAGA 步骤:正向动作 + 对应补偿 + 这一步的业务数据
Store
存储层的统一入口。按 URL 前缀自动选后端:
WorkflowCtx
传给 workflow 函数的上下文。分支都从这里开。

Enums§

Advance
推进全局事务后,状态机给出的下一步指令
Backend
支持的后端
BranchOp
分支操作类型。跟 DTM 的字符串保持一致,方便客户端互通。
BranchResult
分支被调用后的结论。这四态的区分是整个系统的命门。
BranchStatus
Decision
屏障给出的判定
GlobalStatus
TransType
WorkflowError
函数跑不下去的原因。用 ? 往外抛。

Constants§

GRPC_ABORTED
业务明确要求回滚。gRPC 侧的 409
GRPC_FAILED_PRECONDITION
还在处理中。gRPC 侧的 425
GRPC_OK
gRPC 标准状态码。只列用得上的三个,其余一律走 _ => Unknown

Functions§

msg_advance
二阶段消息推进决策。
next_interval
用默认策略退避(10s 起、300s 封顶)
next_interval_with
指数退避:initial → ×2 → … → 封顶 max
saga_advance
SAGA 推进决策。不碰 I/O,所以可以穷举测试。
tcc_advance
TCC 推进决策。
workflow_advance
workflow 推进决策。
xa_advance
XA 推进决策。

Type Aliases§

WorkflowResult