You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 22:01:04