61 lines
1.1 KiB
Go
61 lines
1.1 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++ {
|
|
_, err = producer.Send(context.Background(), &pulsar.ProducerMessage{
|
|
Payload: []byte(data),
|
|
})
|
|
if err == nil {
|
|
break
|
|
}
|
|
closeClient()
|
|
create()
|
|
}
|
|
return err
|
|
}
|
|
|
|
func ClosePulsar() {
|
|
producer.Close()
|
|
client.Close()
|
|
}
|