-
Notifications
You must be signed in to change notification settings - Fork 0
/
producer.go
executable file
·46 lines (39 loc) · 1007 Bytes
/
producer.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
package kafka
import (
"fmt"
"os"
ckafka "github.com/confluentinc/confluent-kafka-go/kafka"
)
func NewKafkaProducer() *ckafka.Producer {
configMap := &ckafka.ConfigMap{
"bootstrap.servers": os.Getenv("kafkaBootstrapServers"),
}
producer, err := ckafka.NewProducer(configMap)
if err != nil {
panic(err)
}
return producer
}
func Publish(msg string, topic string, producer *ckafka.Producer, deliveryChan chan ckafka.Event) error {
message := &ckafka.Message{
TopicPartition: ckafka.TopicPartition{Topic: &topic, Partition: ckafka.PartitionAny},
Value: []byte(msg),
}
err := producer.Produce(message, deliveryChan)
if err != nil {
return err
}
return nil
}
func DeliveryReport(deliveryChan chan ckafka.Event) {
for e := range deliveryChan {
switch ev := e.(type) {
case *ckafka.Message:
if ev.TopicPartition.Error != nil {
fmt.Println("Delivery failed", ev.TopicPartition)
} else {
fmt.Println("Delivery message to:", ev.TopicPartition)
}
}
}
}