-
Notifications
You must be signed in to change notification settings - Fork 3
/
producer.go
48 lines (42 loc) · 1.22 KB
/
producer.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
package rocketmq5Kit
import (
"context"
rmq_client "github.com/apache/rocketmq-clients/golang/v5"
"github.com/richelieu-yang/chimera/v2/src/core/sliceKit"
"time"
)
// NewProducer
/*
PS: In most case, you don't need to create many producers, singletion pattern is more recommended.
@param clientLogPath 客户端日志(blank则输出到控制台)
*/
func NewProducer() (rmq_client.Producer, error) {
if config == nil {
return nil, NotSetupError
}
endpoint := sliceKit.Join(config.Endpoints, ";")
producer, err := rmq_client.NewProducer(&rmq_client.Config{
Endpoint: endpoint,
Credentials: config.Credentials,
})
if err != nil {
return nil, err
}
// Start
if err := producer.Start(); err != nil {
return nil, err
}
return producer, nil
}
// SendMessage Deprecated: 仅供参考如何生产消息
func SendMessage(producer rmq_client.Producer, ctx context.Context, topic string, tag *string, body []byte, messageGroup string, deliveryTimestamp time.Time, keys ...string) ([]*rmq_client.SendReceipt, error) {
msg := &rmq_client.Message{
Topic: topic,
Tag: tag,
Body: body,
}
msg.SetKeys(keys...)
msg.SetMessageGroup(messageGroup)
msg.SetDelayTimestamp(deliveryTimestamp)
return producer.Send(ctx, msg)
}