add pulsar client

This commit is contained in:
w1ndyb0y
2021-02-01 15:35:13 +08:00
parent 9a14fb2cef
commit 3ad4c76acd
8 changed files with 246 additions and 5 deletions
+60
View File
@@ -0,0 +1,60 @@
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()
}
+4 -1
View File
@@ -47,7 +47,7 @@ func Listen(address string) {
}
func WriteToClient(data string) error {
if client != nil {
if Tcp {
_, err := client.Write([]byte(data + "\r\n"))
if err != nil {
Log.Error("error in write data to client socket ", zap.Error(err))
@@ -55,6 +55,9 @@ func WriteToClient(data string) error {
}
return err
}
if Pulsar {
return PulsarSend(data)
}
return nil
}
+2
View File
@@ -15,6 +15,8 @@ const (
var buffer string
var exp, _ = regexp.Compile(Expression)
var RawLog *rotatelogs.RotateLogs
var Tcp bool
var Pulsar bool
func Append(data string) bool {
buffer += data + "\n"