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

Go协程处理函数未执行预期次数问题排查

问题描述

我编写了一个处理Go协程的DoStuff函数,预期会启动noOfLoops*len(objs)个goroutine,但实际执行几次后函数就停滞且无法返回。提供的可复现示例程序本该调用1000次Publish函数,却仅执行100次就停止。

原始函数代码

func (s *service) DoStuff(ctx context.Context, noOfLoops int) error {
    errChan := make(chan error)

    wg := sync.WaitGroup{}

    for i := 1; i <= noOfLoops; i++ {

        offset := 100 * i

        objs, err := s.GetObjs(ctx, offset)
        if err != nil {
            return err
        }

        wg.Add(len(objs))
        for _, svc := range objs {
            srvc := svc

            go func() {
                defer wg.Done()
                
                ... DO STUFF

                if err != nil {
                    errChan <- err
                }
            }()
        }
    }
    wg.Wait()
    err = <-errChan

    if err != nil {
        return err
    }
    return nil
}

可复现示例程序

package main

import (
    "context"
    "fmt"
    "sync"
)

type Obj struct {
    ID string
}

const limit = 1000

func main() {
    db := InitDb()
    err := DoStuff(context.Background(), db, 10)
    if err != nil {
        return
    }

    // mimics running server
    for {

    }
}

func DoStuff(ctx context.Context, db []Obj, noOfLoops int) error {

    wg := sync.WaitGroup{}

    for i := 0; i < noOfLoops; i++ {

        offset := limit * i

        objs, err := GetObjs(ctx, db, offset)
        if err != nil {
            return err
        }

        for _, svc := range objs {
            srvc := svc
            wg.Add(1)
            go func() {
                defer wg.Done()

                err := Publish(ctx, srvc)

                if err != nil {
                    fmt.Println(err)
                }
            }()
        }
    }
    wg.Wait()

    return nil
}

func InitDb() []Obj {
    var out []Obj

    for i := 0; i < 50000; i++ {
        out = append(out, Obj{
            ID: fmt.Sprintf("%d", i),
        })
    }

    return out
}

// mimics getting paged objects from mongo db
func GetObjs(ctx context.Context, db []Obj, offset int) ([]Obj, error) {
    var out []Obj
    
    for i := offset; i < len(db)-(len(db)-offset-limit); i++ { // 100 is limit
        out = append(out, db[i])
    }

    return out, nil
}

// mimics publish to kafka
func Publish(ctx context.Context, obj Obj) error {
    fmt.Println("publishing", obj.ID)
    return nil
}

错误分析

1. 原始函数的通道阻塞问题

原始DoStuff里的errChan是无缓冲通道:

errChan := make(chan error)
  • 若协程在wg.Wait()前尝试向errChan发送错误,会因无接收方直接阻塞,导致defer wg.Done()无法执行,最终wg.Wait()永久等待,函数停滞。
  • 若全程无错误,<-errChan会因无数据发送永久阻塞,函数无法返回。

2. 示例程序的分页逻辑错误

GetObjs里的循环条件写得复杂且逻辑错误:

for i := offset; i < len(db)-(len(db)-offset-limit); i++ {

化简后本应是i < offset + limit,但原代码的写法在offset + limit > len(db)时会出现计算异常,导致实际取出的对象数量远小于预期(比如你遇到的仅100次执行)。


修复方案

针对原始函数的修复

推荐用sync.Once记录第一个错误(避免通道阻塞问题),或者给errChan设置足够的缓冲:

func (s *service) DoStuff(ctx context.Context, noOfLoops int) error {
    var firstErr error
    var once sync.Once // 确保只记录第一个错误
    wg := sync.WaitGroup{}

    for i := 1; i <= noOfLoops; i++ {
        offset := 100 * i
        objs, err := s.GetObjs(ctx, offset)
        if err != nil {
            return err
        }

        wg.Add(len(objs))
        for _, svc := range objs {
            srvc := svc
            go func() {
                defer wg.Done()
                
                // 替换为实际业务逻辑
                err := doActualWork(ctx, srvc)

                if err != nil {
                    // 仅保存第一个出现的错误
                    once.Do(func() {
                        firstErr = err
                    })
                }
            }()
        }
    }
    wg.Wait()

    return firstErr
}

针对示例程序的修复

修正GetObjs的分页逻辑,简化为清晰的范围判断:

func GetObjs(ctx context.Context, db []Obj, offset int) ([]Obj, error) {
    var out []Obj
    end := offset + limit
    // 避免超出数组长度
    if end > len(db) {
        end = len(db)
    }
    // 从offset到end-1遍历取数
    for i := offset; i < end; i++ {
        out = append(out, db[i])
    }
    return out, nil
}

修正后每次能正确取出limit个对象(或剩余的所有对象),确保协程数量符合预期。


内容的提问来源于stack exchange,提问作者user3353167

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 18:47:17