收发普通消息(1) 单向发送 发送方只负责发送消息,不等待服务端返回响应且没有回调函数触发,即只发送请求不等待应答。此方式发送消息的过程耗时非常短,一般在微秒级别。适用于某些耗时非常短,但对可靠性要求并不高的场景,例如日志收集。 参考如下示例代码。 package main import ( "context" "fmt" "github.com/apache/rocketmqclientgo/v2" "github.com/apache/rocketmqclientgo/v2/primitive" "github.com/apache/rocketmqclientgo/v2/producer" "os" "strconv" ) func main() { // 填写分布式消息服务RocketMQ控制台Namesrv接入点 endpoint : "${ENDPOINT}" // 填写AccessKey,在分布式消息服务RocketMQ控制台用户管理菜单中创建的用户ID accessKey : "${ACCESSKEY}" // 填写SecretKey 在分布式消息服务RocketMQ控制台用户管理菜单中创建的用户密钥 secretKey : "${SECRETKEY}" // 填写Topic,在管理控制台创建 topic : "${TOPIC}" p, : rocketmq.NewProducer( producer.WithNsResolver(primitive.NewPassthroughResolver([]string{endpoint})), producer.WithRetry(2), producer.WithCredentials(primitive.Credentials{ AccessKey: accessKey, SecretKey: secretKey, }), ) err : p.Start() if err ! nil { fmt.Printf("start producer error: %s", err.Error()) os.Exit(1) } for i : 0; i < 4; i++ { msg : &primitive.Message{ Topic: topic, Body: []byte("Hello RocketMQ! " + strconv.Itoa(i)), } // 使用单向方式发送消息 err : p.SendOneWay(context.Background(), msg) if err ! nil { fmt.Printf("send message error: %sn", err) } else { fmt.Printf("send message success") } } err p.Shutdown() if err ! nil { fmt.Printf("shutdown producer error: %s", err.Error()) } }
来自: