У меня есть подпрограмма go, которая в основном действует как KafkaConsumer
, она читает сообщения из топи c, а затем порождает еще одну go routine
для каждого полученного сообщения. Теперь этот Consumer go routine
должен быть выключен, когда приложение, которое является main go routine
, закрывается. Но я сталкиваюсь с трудностями в правильном закрытии этого. Ниже приведено определение Kafka Consumer
package svc
import (
"event-service/pkg/pb"
"fmt"
"github.com/gogo/protobuf/proto"
"gopkg.in/confluentinc/confluent-kafka-go.v1/kafka"
"log"
"os"
"sync"
)
type EventConsumer func(event eventService.Event)
type KafkaConsumer struct {
done chan bool
eventChannels []string
consumer *kafka.Consumer
consumerMapping map[string]EventConsumer
wg *sync.WaitGroup
}
func getKafkaConsumerConfigMap(config map[string]interface{}) *kafka.ConfigMap {
configMap := &kafka.ConfigMap{}
for key, value := range config {
err := configMap.SetKey(key, value)
if err != nil {
log.Println(fmt.Sprintf("An error %v occurred while setting %v: %v", err, key, value))
}
}
return configMap
}
func NewKafkaConsumer(channels []string, config map[string]interface{}, consumerMapping map[string]EventConsumer) *KafkaConsumer {
var wg sync.WaitGroup
consumer, err := kafka.NewConsumer(getKafkaConsumerConfigMap(config))
done := make(chan bool, 1)
if err != nil {
log.Fatalf("An error %v occurred while starting kafka consumer.", err)
}
err = consumer.SubscribeTopics(channels, nil)
if err != nil {
log.Fatalf("An error %v occurred while subscribing to kafka topics %v.", err, channels)
}
return &KafkaConsumer{eventChannels: channels, done: done, wg: &wg, consumer: consumer, consumerMapping: consumerMapping}
}
func (kc *KafkaConsumer) getEvent(eventData []byte) *eventService.Event {
event := eventService.Event{}
err := proto.Unmarshal(eventData, &event)
if err != nil {
log.Println(fmt.Sprintf("An error %v occurred while un marshalling data from kafka.", err))
}
return &event
}
func (kc *KafkaConsumer) Consume() {
go func() {
run := true
for run == true {
select {
case sig := <-kc.done:
log.Println(fmt.Sprintf("Caught signal %v: terminating \n", sig))
run = false
return
default:
}
e := <-kc.consumer.Events()
switch event := e.(type) {
case kafka.AssignedPartitions:
_, _ = fmt.Fprintf(os.Stderr, "%% %v\n", event)
err := kc.consumer.Assign(event.Partitions)
if err != nil {
log.Println(fmt.Sprintf("An error %v occurred while assigning partitions.", err))
}
case kafka.RevokedPartitions:
_, _ = fmt.Fprintf(os.Stderr, "%% %v\n", event)
err := kc.consumer.Unassign()
if err != nil {
log.Println(fmt.Sprintf("An error %v occurred while unassigning partitions.", err))
}
case *kafka.Message:
domainEvent := kc.getEvent(event.Value)
kc.wg.Add(1)
go func(event *eventService.Event) {
defer kc.wg.Done()
if eventConsumer := kc.consumerMapping[domainEvent.EntityType]; eventConsumer != nil {
eventConsumer(*domainEvent)
} else {
log.Println(fmt.Sprintf("Event consumer not found for %v event type", domainEvent.EntityType))
}
}(domainEvent)
case kafka.PartitionEOF:
fmt.Printf("%% Reached %v\n", e)
case kafka.Error:
_, _ = fmt.Fprintf(os.Stderr, "%% Error: %v\n", e)
}
}
}()
}
func (kc *KafkaConsumer) Close() {
log.Println("Waiting")
kc.wg.Wait()
kc.done <- true
log.Println("Done waiting")
err := kc.consumer.Close()
if err != nil {
log.Println(fmt.Sprintf("An error %v occurred while closing kafka consumer.", err))
}
}
. Ниже приведен код основного потока
package main
import (
"event-service/pkg/pb"
"event-service/pkg/svc"
"fmt"
"log"
)
func main() {
eventConsumerMapping := map[string]svc.EventConsumer{"doctor-created": func(event eventService.Event) {
log.Println(fmt.Sprintf("Got event %v from kafka", event))
}}
consumerConfig := map[string]interface{}{
"bootstrap.servers": "localhost:9092",
"group.id": "catalog",
"go.events.channel.enable": true,
"go.application.rebalance.enable": true,
"enable.partition.eof": true,
"auto.offset.reset": "earliest",
}
kafkaConsumer := svc.NewKafkaConsumer([]string{"doctor-created"}, consumerConfig, eventConsumerMapping)
kafkaConsumer.Consume()
kafkaConsumer.Close()
}
. Проблема в том, что приложение иногда вообще не заканчивается и не выполняет consume
функция в некоторых прогонах, что мне здесь не хватает?