终止Go例程中的无限循环
来源:stackoverflow
时间:2024-02-27 16:03:27 349浏览 收藏
一分耕耘,一分收获!既然打开了这篇文章《终止Go例程中的无限循环》,就坚持看下去吧!文中内容包含等等知识点...希望你能在阅读本文后,能真真实实学到知识或者帮你解决心中的疑惑,也欢迎大佬或者新人朋友们多留言评论,多给建议!谢谢!
问题内容
我有一个 go 例程,它基本上充当 kafkaconsumer
,它从主题读取消息,然后为收到的每条消息生成另一个 go 例程
。现在,当 main goroutine
应用程序关闭时,这个 consumer goroutine
应该关闭。但我在正确关闭它方面遇到了困难。
下面是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
函数,我在这里缺少什么?
解决方案
好的,解决方案如下, 1. 由于只要主 go 例程还活着,消费者 go 例程就应该存在,并且主 go 例程也是一个无尽的 go 例程,因此在 go 例程运行时关闭消费者 go 例程并不是正确的方法。
因此以下解决方案有效
package main import ( "event-service/pkg/pb" "event-service/pkg/svc" "fmt" "log" "sync" ) 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-2", "session.timeout.ms": 6000, "go.events.channel.enable": true, "go.application.rebalance.enable": true, "enable.partition.eof": true, "auto.offset.reset": "earliest", } var wg sync.waitgroup kafkaconsumer := svc.newkafkaconsumer([]string{"doctor-created"}, consumerconfig, eventconsumermapping) kafkaconsumer.consume(&wg) wg.wait() kafkaconsumer.close() }
服务定义
package svc import ( "event-service/pkg/pb" "fmt" "github.com/gogo/protobuf/proto" "gopkg.in/confluentinc/confluent-kafka-go.v1/kafka" "log" "os" "os/signal" "sync" "syscall" ) type EventConsumer func(event eventService.Event) type KafkaConsumer struct { done chan bool consumer *kafka.Consumer consumerMapping map[string]EventConsumer sigChan chan os.Signal channels []string } 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 { sigChan := make(chan os.Signal, 1) 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{channels: channels, sigChan: sigChan, done: done, 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(wg *sync.WaitGroup) { signal.Notify(kc.sigChan, syscall.SIGINT, syscall.SIGTERM) wg.Add(1) go func() { run := true defer wg.Done() for run == true { select { case sig := <-kc.sigChan: fmt.Printf("Caught signal %v: terminating\n", sig) run = false case ev := <-kc.consumer.Events(): switch e := ev.(type) { case kafka.AssignedPartitions: _, _ = fmt.Fprintf(os.Stderr, "%% %v\n", e) _ = kc.consumer.Assign(e.Partitions) case kafka.RevokedPartitions: _, _ = fmt.Fprintf(os.Stderr, "%% %v\n", e) _ = kc.consumer.Unassign() case *kafka.Message: domainEvent := kc.getEvent(e.Value) wg.Add(1) go func(event *eventService.Event) { defer 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: // Errors should generally be considered as informational, the client will try to automatically recover _, _ = fmt.Fprintf(os.Stderr, "%% Error: %v\n", e) } } } }() } func (kc *KafkaConsumer) Close() { err := kc.consumer.Close() if err != nil { log.Println(fmt.Sprintf("An error %v occurred while closing kafka consumer.", err)) } }
以上就是《终止Go例程中的无限循环》的详细内容,更多关于的资料请关注golang学习网公众号!
声明:本文转载于:stackoverflow 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
-
502 收藏
-
502 收藏
-
501 收藏
-
501 收藏
-
501 收藏
最新阅读
更多>
-
139 收藏
-
204 收藏
-
325 收藏
-
477 收藏
-
486 收藏
-
439 收藏
-
357 收藏
-
352 收藏
-
101 收藏
-
440 收藏
-
212 收藏
-
143 收藏
课程推荐
更多>
-
- 前端进阶之JavaScript设计模式
- 设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
- 立即学习 542次学习
-
- GO语言核心编程课程
- 本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
- 立即学习 508次学习
-
- 简单聊聊mysql8与网络通信
- 如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
- 立即学习 497次学习
-
- JavaScript正则表达式基础与实战
- 在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
- 立即学习 487次学习
-
- 从零制作响应式网站—Grid布局
- 本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
- 立即学习 484次学习