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 }