10 changed files with 187 additions and 60 deletions
@ -1,17 +1,34 @@ |
|||||
### 轻量级分布式事务管理服务 |
### 轻量级分布式事务管理服务 |
||||
|
跨语言--语言无关,基于http协议 |
||||
## 配置rabbitmq和mysql |
支持xa、tcc、saga |
||||
|
## 快速开始 |
||||
dtm依赖于rabbitmq和mysql,请搭建好rabbitmq和mysql,并修改dtm.yml |
场景描述: |
||||
|
假设您实现了一个转账功能,分为两个微服务:转入、转出 |
||||
## 启动tc |
转出:服务地址为 http://example.com/api/busi_saga/transOut?gid=xxx POST 参数为 {"uid": 2, "amount":30} |
||||
|
转入:服务地址为 http://example.com/api/busi_saga/transIn?gid=xxx POST 参数为 {"uid": 1, "amount":30} |
||||
```go run dtm-svr/svr``` |
在saga模式下,有对应的补偿微服务 |
||||
|
转出:服务地址为 http://example.com/api/busi_saga/transOutCompensate?gid=xxx POST 参数为 {"uid": 2, "amount":30} |
||||
## 启动例子saga的tm+rm |
转入:服务地址为 http://example.com/api/busi_saga/transInCompensate?gid=xxx POST 参数为 {"uid": 1, "amount":30} |
||||
|
HTTP协议方式 |
||||
```go run example/saga``` |
curl -d '{"gid":"xxx","trans_type":"saga","steps":[{"action":"http://example.com/api/busi_saga/TransOut","compensate":"http://example.com/api/busi_saga/TransOutCompensate","data":"{\"amount\":30}"},{"action":"http://localhost:8081/api/busi_saga/TransIn","compensate":"http://localhost:8081/api/busi_saga/TransInCompensate","data":"{\"amount\":30}"}]}' 8.140.124.252/api/dtm/commit |
||||
|
此请求向dtm提交了一个saga事务,dtm会按照saga模式,请求transIn/transOut,并且在出错情况下,保证抵用相关的补偿api |
||||
## 或者启动例子tcc的tm+rm |
go客户端方式 |
||||
|
// 事务参与者的服务地址 |
||||
```go run example/tcc``` |
const startBusiPort = 8084 |
||||
|
const startBusiApi = "/api/busi_start" |
||||
|
|
||||
|
var startBusi = fmt.Sprintf("http://localhost:%d%s", startBusiPort, startBusiApi) |
||||
|
err := dtm.SagaNew(DtmServer, gid).Add(startBusi+"/TransOut", startBusi+"/TransOutCompensate", &gin.H{ |
||||
|
"amount": 30, |
||||
|
"uid": 2, |
||||
|
}).Add(startBusi+"/TransIn", startBusi+"/TransInCompensate", &gin.H{ |
||||
|
"amount": 30, |
||||
|
"uid": 1 |
||||
|
}).Commit() |
||||
|
|
||||
|
本地启动方式 |
||||
|
需要安装docker,和docker-compose |
||||
|
curl localhost:8080/api/initMysql |
||||
|
go run examples/app/main saga |
||||
|
|
||||
|
其他 |
||||
|
|||||
@ -0,0 +1,76 @@ |
|||||
|
package examples |
||||
|
|
||||
|
import ( |
||||
|
"fmt" |
||||
|
"time" |
||||
|
|
||||
|
"github.com/gin-gonic/gin" |
||||
|
"github.com/sirupsen/logrus" |
||||
|
"github.com/yedf/dtm" |
||||
|
"github.com/yedf/dtm/common" |
||||
|
) |
||||
|
|
||||
|
// 事务参与者的服务地址
|
||||
|
const startBusiPort = 8084 |
||||
|
const startBusiApi = "/api/busi_start" |
||||
|
|
||||
|
var startBusi = fmt.Sprintf("http://localhost:%d%s", startBusiPort, startBusiApi) |
||||
|
|
||||
|
func startMain() { |
||||
|
go startStartSvr() |
||||
|
startFireRequest() |
||||
|
time.Sleep(1000 * time.Second) |
||||
|
} |
||||
|
|
||||
|
func startStartSvr() { |
||||
|
logrus.Printf("saga examples starting") |
||||
|
app := common.GetGinApp() |
||||
|
startAddRoute(app) |
||||
|
app.Run(fmt.Sprintf(":%d", SagaBusiPort)) |
||||
|
} |
||||
|
|
||||
|
func startFireRequest() { |
||||
|
gid := common.GenGid() |
||||
|
logrus.Printf("busi transaction begin: %s", gid) |
||||
|
req := &TransReq{ |
||||
|
Amount: 30, |
||||
|
TransInResult: "SUCCESS", |
||||
|
TransOutResult: "SUCCESS", |
||||
|
} |
||||
|
saga := dtm.SagaNew(DtmServer, gid). |
||||
|
Add(startBusi+"/TransOut", startBusi+"/TransOutCompensate", req). |
||||
|
Add(startBusi+"/TransIn", startBusi+"/TransInCompensate", req) |
||||
|
logrus.Printf("busi trans commit") |
||||
|
err := saga.Commit() |
||||
|
e2p(err) |
||||
|
} |
||||
|
|
||||
|
func startAddRoute(app *gin.Engine) { |
||||
|
app.POST(SagaBusiApi+"/TransIn", common.WrapHandler(startTransIn)) |
||||
|
app.POST(SagaBusiApi+"/TransInCompensate", common.WrapHandler(startTransInCompensate)) |
||||
|
app.POST(SagaBusiApi+"/TransOut", common.WrapHandler(startTransOut)) |
||||
|
app.POST(SagaBusiApi+"/TransOutCompensate", common.WrapHandler(startTransOutCompensate)) |
||||
|
logrus.Printf("examples listening at %d", startBusiPort) |
||||
|
} |
||||
|
|
||||
|
func startTransIn(c *gin.Context) (interface{}, error) { |
||||
|
gid := c.Query("gid") |
||||
|
req := transReqFromContext(c) |
||||
|
logrus.Printf("%s TransIn: %v result: %s", gid, req, req.TransInResult) |
||||
|
return M{"result": req.TransInResult}, nil |
||||
|
} |
||||
|
|
||||
|
func startTransInCompensate(c *gin.Context) (interface{}, error) { |
||||
|
return M{"result": "SUCCESS"}, nil |
||||
|
} |
||||
|
|
||||
|
func startTransOut(c *gin.Context) (interface{}, error) { |
||||
|
gid := c.Query("gid") |
||||
|
req := transReqFromContext(c) |
||||
|
logrus.Printf("%s TransOut: %v result: %s", gid, req, req.TransOutResult) |
||||
|
return M{"result": req.TransOutResult}, nil |
||||
|
} |
||||
|
|
||||
|
func startTransOutCompensate(c *gin.Context) (interface{}, error) { |
||||
|
return M{"result": "SUCCESS"}, nil |
||||
|
} |
||||
Loading…
Reference in new issue