package datacenter import ( "91porn-server/common" "91porn-server/common/log" "time" "github.com/IBM/sarama" ) // 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 } // AsyncSendMessage 同步发送确保消息成功 func AsyncSendMessage(topic TopicType, message []byte) { common.Go(func() { if gSyncProducer != nil { msg := sarama.ProducerMessage{ Topic: string(topic), Value: sarama.ByteEncoder(message), } partition, offset, err := gSyncProducer.SendMessage(&msg) if err != nil { log.Error("gSyncProducer SendMessage Fail", log.Any("topic", topic), log.E(err)) return } log.Info("SyncSendMessage ", log.Any("topic", topic), log.Any("Partition", partition), log.Any("Offset", offset)) } }) } // 创建同步生产者 用于对消息的顺序有严格要求的场景 性能相对较低 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 }