Files
rootandClaude Opus 5 8679200f41 Initial commit
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-15 13:57:10 +08:00

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
}