feat kafka
This commit is contained in:
@@ -9,12 +9,11 @@ import (
|
||||
var client *Client
|
||||
|
||||
type Client struct {
|
||||
producer sarama.AsyncProducer
|
||||
consumer sarama.ConsumerGroup
|
||||
serverName string
|
||||
producer sarama.AsyncProducer
|
||||
consumer sarama.ConsumerGroup
|
||||
}
|
||||
|
||||
func Init(cfg *config.KafkaConfig, serverName string) error {
|
||||
func Init(cfg *config.KafkaConfig) error {
|
||||
producer, err := getAsyncProducer(cfg)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -24,9 +23,8 @@ func Init(cfg *config.KafkaConfig, serverName string) error {
|
||||
return err
|
||||
}
|
||||
client = &Client{
|
||||
producer: producer,
|
||||
consumer: consumer,
|
||||
serverName: serverName,
|
||||
producer: producer,
|
||||
consumer: consumer,
|
||||
}
|
||||
go producerError()
|
||||
go consumerError()
|
||||
|
||||
Reference in New Issue
Block a user