75 lines
2.0 KiB
Go
75 lines
2.0 KiB
Go
package kafka
|
|
|
|
import (
|
|
"fmt"
|
|
"time"
|
|
|
|
"91porn-server/common/log"
|
|
|
|
"github.com/Shopify/sarama"
|
|
)
|
|
|
|
const ClientID string = "pf_sp_kafka"
|
|
|
|
// kafka 使用
|
|
var gSyncProducer sarama.SyncProducer
|
|
|
|
/**********************************************客户端 *********************/
|
|
// 初始化kafka生产者 发送消息入口
|
|
func InitKafkaProducter(addrs []string) error {
|
|
config := sarama.NewConfig()
|
|
config.Version = sarama.V2_0_0_0
|
|
config.Producer.Return.Successes = true
|
|
config.Net.KeepAlive = 2 * time.Hour
|
|
cli, err := sarama.NewClient(addrs, config)
|
|
if err != nil {
|
|
log.Error("startUp Kafka Init Kafka error", log.E(err))
|
|
return err
|
|
}
|
|
gSyncProducer, err = SyncProducter(cli)
|
|
if err != nil {
|
|
log.Error("startUp Kafka new SyncProducter error", log.E(err))
|
|
return nil
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SyncSendMessage 同步发送确保消息成功
|
|
func SyncSendMessage(topic string, message []byte) {
|
|
if gSyncProducer != nil {
|
|
msg := sarama.ProducerMessage{
|
|
Topic: topic,
|
|
Value: sarama.ByteEncoder(message),
|
|
}
|
|
if _, _, err := gSyncProducer.SendMessage(&msg); err != nil {
|
|
log.Error("gSyncProducer SendMessage Fail", log.Any("topic", topic), log.E(err))
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
// 创建同步生产者 用于对消息的顺序有严格要求的场景 性能相对较低
|
|
func SyncProducter(client sarama.Client) (sarama.SyncProducer, error) {
|
|
producer, err := sarama.NewSyncProducerFromClient(client)
|
|
if err != nil {
|
|
log.Error("create syncProducer error", log.E(err))
|
|
return nil, err
|
|
}
|
|
return producer, nil
|
|
}
|
|
|
|
/**********************************************客户端 end*********************/
|
|
|
|
/**********************************************服务端 *********************/
|
|
func InitKafkaConsumerGroup(addrs []string, groupName string) sarama.ConsumerGroup {
|
|
config := sarama.NewConfig()
|
|
config.Version = sarama.V2_0_0_0
|
|
config.Consumer.Return.Errors = true
|
|
fmt.Println("addrs:", addrs)
|
|
group, err := sarama.NewConsumerGroup(addrs, groupName, config)
|
|
if err != nil {
|
|
panic(err)
|
|
}
|
|
return group
|
|
}
|