You need to enable JavaScript to run this app.
文档中心
消息队列 RocketMQ版

消息队列 RocketMQ版

复制全文
下载 pdf
Go 5.x SDK
事务消息
复制全文
下载 pdf
事务消息
本文提供使用 Go SDK 收发事务消息的示例代码供您参考。
发送事务消息
发送事务消息需要在控制台创建事务消息类型的 Topic,在发送 half 消息后,用户可以根据自身业务逻辑对消息选择:
  • Transaction.COMMIT:提交事务。事务提交之后,之前暂存的消息对消费者可见。
  • Transaction.ROLLBACK:回滚事务。事务回滚,之前暂存的消息会被标记为已经回滚。
此外,事务消息需要自定义checker,以便服务端在未收到确认时触发回查。
package main
import (
"context"
"fmt"
"log"
"os"
"strconv"
"time"
rmq_client "github.com/apache/rocketmq-clients/golang/v5"
"github.com/apache/rocketmq-clients/golang/v5/credentials"
)
// 设置为您从火山引擎消息队列 RocketMQ版控制台获取的接入点信息,类似"{INSTANCE_ID}.rocketmq.ivolces.com:8080"。
// 根据实例支持的认证方式,填写对应的身份认证信息:
// 密钥管理:AccessKey 为 AccessKey ID,SecretKey 为 AccessKey Secret。
// 权限管理:AccessKey 为 ACL 用户名,SecretKey 为 ACL 密码。
const (
Topic = "xxxx"
Endpoint = "rocketmq-xxxx:8080"
AccessKey = "xxxx"
SecretKey = "xxxx"
)
func main() {
// log to console
os.Setenv("mq.consoleAppender.enabled", "true")
rmq_client.ResetLogger()
// In most case, you don't need to create many producers, singleton pattern is more recommended.
producer, err := rmq_client.NewProducer(&rmq_client.Config{
Endpoint: Endpoint,
Credentials: &credentials.SessionCredentials{
AccessKey: AccessKey,
AccessSecret: SecretKey,
},
},
rmq_client.WithTransactionChecker(&rmq_client.TransactionChecker{
Check: func(msg *rmq_client.MessageView) rmq_client.TransactionResolution {
log.Printf("check transaction message: %v", msg)
return rmq_client.COMMIT
},
}),
rmq_client.WithTopics(Topic),
)
if err != nil {
log.Fatal(err)
}
// start producer
err = producer.Start()
if err != nil {
log.Fatal(err)
}
// graceful stop producer
defer producer.GracefulStop()
for i := 0; i < 10; i++ {
// new a message
msg := &rmq_client.Message{
Topic: Topic,
Body: []byte("this is a message : " + strconv.Itoa(i)),
}
// set keys and tag
msg.SetKeys("a", "b")
msg.SetTag("ab")
// send message in sync
transaction := producer.BeginTransaction()
resp, err := producer.SendWithTransaction(context.TODO(), msg, transaction)
if err != nil {
log.Fatal(err)
}
for i := 0; i < len(resp); i++ {
fmt.Printf("%#v\n", resp[i])
}
// commit transaction message
err = transaction.Commit()
if err != nil {
log.Fatal(err)
}
// wait a moment
time.Sleep(time.Second * 1)
}
}
订阅事务消息
事务消息的订阅方式与普通消息一致,示例代码如下所示。
package main
import (
"context"
"fmt"
"log"
"os"
"time"
rmq_client "github.com/apache/rocketmq-clients/golang/v5"
"github.com/apache/rocketmq-clients/golang/v5/credentials"
)
// 设置为您从火山引擎消息队列 RocketMQ版控制台获取的接入点信息,类似"{INSTANCE_ID}.rocketmq.ivolces.com:8080"。
// 根据实例支持的认证方式,填写对应的身份认证信息:
// 密钥管理:AccessKey 为 AccessKey ID,SecretKey 为 AccessKey Secret。
// 权限管理:AccessKey 为 ACL 用户名,SecretKey 为 ACL 密码。
const (
Topic = "xxxx"
ConsumerGroup = "xxxx"
Endpoint = "rocketmq-xxxx:8080"
AccessKey = "xxxx"
SecretKey = "xxxx"
)
var (
// maximum waiting time for receive func
awaitDuration = time.Second * 5
// maximum number of messages received at one time
maxMessageNum int32 = 16
// invisibleDuration should > 20s
invisibleDuration = time.Second * 20
// receive messages in a loop
receiveConcurrency = 4
)
func main() {
// log to console
os.Setenv("mq.consoleAppender.enabled", "true")
rmq_client.ResetLogger()
// In most case, you don't need to create many consumers, singleton pattern is more recommended.
simpleConsumer, err := rmq_client.NewSimpleConsumer(&rmq_client.Config{
Endpoint: Endpoint,
ConsumerGroup: ConsumerGroup,
Credentials: &credentials.SessionCredentials{
AccessKey: AccessKey,
AccessSecret: SecretKey,
},
},
rmq_client.WithAwaitDuration(awaitDuration),
rmq_client.WithSubscriptionExpressions(map[string]*rmq_client.FilterExpression{
Topic: rmq_client.SUB_ALL,
}),
)
if err != nil {
log.Fatal(err)
}
// start simpleConsumer
err = simpleConsumer.Start()
if err != nil {
log.Fatal(err)
}
// graceful stop simpleConsumer
defer simpleConsumer.GracefulStop()
go func() {
for {
fmt.Println("start receive message")
mvs, err := simpleConsumer.Receive(context.TODO(), maxMessageNum, invisibleDuration)
if err != nil {
fmt.Println(err)
}
// ack message
for _, mv := range mvs {
simpleConsumer.Ack(context.TODO(), mv)
fmt.Println(mv)
}
fmt.Println("wait a moment")
fmt.Println()
}
}()
// run for a while
time.Sleep(100 * time.Minute)
}
最近更新时间:2026.08.26 15:19:56
这个页面对您有帮助吗?
有用
有用
无用
无用