Goroutine实现Kafka消费者工作组扩容问题求助
解决方案:Kafka多消费者工作组的正确实现
问题根源
原代码中同步调用StartWorker,当工作组规模大于1时,第一个Worker会阻塞主线程,后续Worker无法启动;同时多个Worker共用一个信号通道,信号被第一个Worker读取后,其他Worker无法感知退出信号,导致整体无法正常工作。
改造方案
通过context.WithCancel统一管理退出信号,结合sync.WaitGroup等待所有Worker协程完成,同时将每个Worker的启动改为异步协程:
1. 改造main函数
引入context和sync包,创建可取消的Context,使用WaitGroup等待所有Worker退出:
package main import ( "context" "db_write_consumer/db" "db_write_consumer/worker" "fmt" "os" "os/signal" "sync" "syscall" ) func main() { // 创建可取消的Context,用于统一控制所有Worker退出 ctx, cancel := context.WithCancel(context.Background()) defer cancel() // 监听系统退出信号,异步触发Context取消 sigchan := make(chan os.Signal, 1) signal.Notify(sigchan, syscall.SIGINT, syscall.SIGTERM) go func() { <-sigchan fmt.Println("收到退出信号,开始终止所有Worker...") cancel() }() // 初始化MySQL客户端,完善错误处理 mySQLClient, err := db.NewMySQLDBClient("root", "", "localhost", 3306, "testbase") if err != nil { fmt.Fprintf(os.Stderr, "初始化MySQL客户端失败: %v\n", err) os.Exit(1) } defer mySQLClient.Close() // 确保程序退出时关闭数据库连接 // 创建消费者工作组(可设置大于1的数量) workerCount := 3 workers := worker.CreateGroup("localhost:9092", "testgroup", workerCount) // 使用WaitGroup等待所有Worker协程完成 var wg sync.WaitGroup wg.Add(workerCount) for _, w := range workers { worker := w go func() { defer wg.Done() // 协程结束时通知WaitGroup worker.StartWorker(ctx, []string{"test-topic"}, mySQLClient) }() } // 阻塞主线程,等待所有Worker退出 wg.Wait() fmt.Println("所有Worker已终止,程序退出") }
2. 改造StartWorker函数
将原来的sigchan替换为context.Context,监听ctx.Done()信号终止循环,同时完善错误处理:
import ( "context" "fmt" "os" "database/sql" "github.com/confluentinc/confluent-kafka-go/kafka" "google.golang.org/protobuf/proto" "your-project-path/pb" // 替换为实际的proto包路径 ) func StartWorker(ctx context.Context, c *kafka.Consumer, topics []string, mySQLClient *sql.DB) { // 订阅主题,处理订阅错误 err := c.SubscribeTopics(topics, nil) if err != nil { fmt.Fprintf(os.Stderr, "订阅主题失败: %v\n", err) return } defer func() { fmt.Printf("关闭消费者: %v\n", c) c.Close() }() fmt.Printf("启动消费者: %v\n", c) for { select { case <-ctx.Done(): // 监听Context取消信号,触发退出 fmt.Println("收到终止信号,停止消费") return default: // 读取消息,处理超时和读取错误 ev, err := c.ReadMessage(100) if err != nil { if err.(kafka.Error).Code() == kafka.ErrTimedOut { continue // 超时忽略,继续循环 } fmt.Fprintf(os.Stderr, "读取消息失败: %v\n", err) continue } if ev == nil { continue } // 反序列化Proto消息 msg := &pb.Person{} err = proto.Unmarshal(ev.Value, msg) if err != nil { fmt.Fprintf(os.Stderr, "反序列化消息失败: %v\n", err) continue } // 写入数据库,确保WriteStuff线程安全 err = WriteStuff(mySQLClient, msg.Id, msg.Lastname, msg.Firstname, msg.Address, msg.City) if err != nil { fmt.Fprintf(os.Stderr, "写入数据库失败: %v\n", err) continue } // 打印消息头(可选) if ev.Headers != nil { fmt.Printf("消息头: %v\n", ev.Headers) } // 提交偏移量,处理提交错误 _, err = c.StoreMessage(ev) if err != nil { fmt.Fprintf(os.Stderr, "提交偏移量失败: %v\n", ev.TopicPartition) } } } }
3. 关键注意事项
- MySQL客户端并发安全:
sql.DB本身是并发安全的,可直接在多协程中共享;若WriteStuff有自定义事务或状态操作,需通过sync.Mutex保证线程安全。 - 消费者组配置:确保所有消费者使用相同的
group.id,auto.offset.reset等参数符合业务需求。 - 资源清理:通过
defer确保消费者、数据库连接在退出时被正确关闭,避免资源泄漏。 - 错误处理:新增必要的错误打印,便于排查消费、序列化、数据库写入等环节的问题。
内容的提问来源于stack exchange,提问作者ILoveCliques
相关产品推荐
相关产品推荐

