Files
tele-recv/utils/pulsar.go
T
2021-02-01 17:46:20 +08:00

65 lines
1.2 KiB
Go

package utils
import (
"context"
"github.com/apache/pulsar-client-go/pulsar"
"go.uber.org/zap"
)
var (
pulsarClient pulsar.Client
producer pulsar.Producer
pulsarUrl string
pulsarTopic string
clientName string
)
func CreateProducer(url string, topic string, name string) {
pulsarUrl = url
pulsarTopic = topic
clientName = name
create()
}
func create() {
var err error
pulsarClient, err = pulsar.NewClient(pulsar.ClientOptions{
URL: pulsarUrl,
})
if err != nil {
Log.Error("error in create pulsar client ", zap.Error(err))
}
producer, err = pulsarClient.CreateProducer(pulsar.ProducerOptions{
Topic: pulsarTopic,
Name: clientName,
})
if err != nil {
Log.Error("error in create pulsar producer ", zap.Error(err))
}
}
func PulsarSend(data string) error {
var err error
for i := 0; i < 5; i++ {
id, err := producer.Send(context.Background(), &pulsar.ProducerMessage{
Payload: []byte(data),
})
if err == nil {
Log.Info("send to pulsar", zap.Any("id", id))
break
}
Log.Error("error in produce message", zap.Error(err))
closeClient()
create()
}
return err
}
func ClosePulsar() {
producer.Close()
if client != nil {
client.Close()
}
}