-
Notifications
You must be signed in to change notification settings - Fork 2
/
Copy pathcontainer.go
59 lines (53 loc) · 1.49 KB
/
container.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
package kafkaclient
import (
"github.com/confluentinc/confluent-kafka-go/kafka"
)
// Container contains all variables and configs to run kafka.
type Container struct {
configMap kafka.ConfigMap
closeSignal chan bool
resourcesCounter uint64
}
// NewContainer initialize a new container.
func NewContainer(config kafka.ConfigMap) *Container {
if config == nil {
config = kafka.ConfigMap{}
}
return &Container{
configMap: config,
closeSignal: make(chan bool),
}
}
// NewProducer initialize a new producer from this container.
func (c *Container) NewProducer(config kafka.ConfigMap) (prod *Producer, err error) {
newConfig := c.initConfig(config)
originProducer, err := kafka.NewProducer(&newConfig)
if err != nil {
return
}
prod = &Producer{originProducer}
c.close(prod)
return
}
// NewAdminClient initialize a new admin client from this container.
func (c *Container) NewAdminClient(config kafka.ConfigMap) (ac *AdminClient, err error) {
newConfig := c.initConfig(config)
originAC, err := kafka.NewAdminClient(&newConfig)
if err != nil {
return
}
ac = &AdminClient{originAC}
c.close(ac)
return
}
// NewConsumer initialize a new consumer from this container.
func (c *Container) NewConsumer(config kafka.ConfigMap) (cons *Consumer, err error) {
newConfig := c.initConfig(config)
originCons, err := kafka.NewConsumer(&newConfig)
if err != nil {
return
}
cons = &Consumer{originCons}
c.close(cons)
return
}