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

Go语言sync.WaitGroup未生效 部分goroutine未执行完成就退出

Go sync.WaitGroup不生效问题排查与修复

根因分析

你的代码存在两个核心逻辑错误,直接导致WaitGroup未达预期效果:

  1. 通道收发次数不匹配:每个IntradayStratify goroutine处理单个标的时,会遍历所有时间周期(从代码看是5/15/30/60四个周期),每个周期无论是否匹配到策略都会往通道发送1条数据,但是你在主循环遍历标的时,每个标的仅从通道读取1次就处理下一个标的,剩余的发送操作直接阻塞在goroutine中,大量goroutine无法执行到defer wg.Done()逻辑。
  2. 主逻辑未等待任务执行完成:你单独启动了一个goroutine执行wg.Wait()和关闭通道的逻辑,但主循环完全没有监听通道关闭信号,遍历完所有标的后直接跳过等待步骤,执行Discord推送逻辑,此时大量策略匹配结果还未被读取处理。

修复方案

1. 调整主逻辑的消费与等待流程

将通道消费逻辑独立为单独的goroutine,主逻辑先等待所有worker执行完成、关闭通道后再执行推送:

func RunIntradayScanner() {
    var wg sync.WaitGroup
    var consumeWg sync.WaitGroup
    logrus.Info("Clearing out pattern slices...")
    var tf5 []request.StratNotification
    var tf15 []request.StratNotification
    var tf30 []request.StratNotification
    var tf60 []request.StratNotification

    var intradayChannel = make(chan request.StratNotification)
    symbols := sources.GetSymbols()
    wg.Add(len(symbols))

    // 启动消费协程,持续读通道直到关闭
    consumeWg.Add(1)
    go func() {
        defer consumeWg.Done()
        for match := range intradayChannel {
            // 跳过空通知
            if match.TimeFrame == 0 {
                continue
            }
            switch match.TimeFrame {
            case 5:
                tf5 = append(tf5, match)
            case 15:
                tf15 = append(tf15, match)
            case 30:
                tf30 = append(tf30, match)
            case 60:
                tf60 = append(tf60, match)
            }
        }
    }()

    // 启动所有扫描worker
    for _, s := range symbols {
        go IntradayStratify(strings.TrimSpace(s.Symbol), intradayChannel, &wg)
    }

    // 等待所有worker执行完成
    logrus.Info("------Waiting for workers to finish")
    wg.Wait()
    logrus.Info("------Closing intraday channel")
    close(intradayChannel)
    
    // 等待所有通道数据消费完成
    consumeWg.Wait()

    // 执行推送逻辑
    if len(tf5) > 0 {
        SplitUpAndSendEmbedToDiscord(5, tf5)
    }
    if len(tf15) > 0 {
        SplitUpAndSendEmbedToDiscord(15, tf15)
    }
    if len(tf30) > 0 {
        SplitUpAndSendEmbedToDiscord(30, tf30)
    }
    if len(tf60) > 0 {
        SplitUpAndSendEmbedToDiscord(60, tf60)
    }
}

2. 优化worker的通道发送逻辑

移除不必要的空数据发送,减少通道开销:

// IntradayStratify - go routine to run during market hours
func IntradayStratify(ticker string, c chan request.StratNotification, wg *sync.WaitGroup) {
    defer wg.Done()

    candles := request.GetIntraday(ticker)
    for _, tf := range timeframes {
        chunkedCandles := request.DetermineTimeframes(tf, ticker, candles)
        if len(chunkedCandles) > 1 {
            highLows := request.CalculateIntraDayHighLow(chunkedCandles)
            if len(highLows) > 2 {
                bl, stratPattern := request.DetermineStratPattern(ticker, tf, highLows)
                if bl {
                    c <- stratPattern
                }
            }
        }
        // 删掉无用的空数据发送逻辑
        // c <- request.StratNotification{}
    }
}

修复效果

所有worker执行完成后才会关闭通道,所有匹配到的策略都会被写入对应切片,不会出现匹配日志存在但无对应推送记录的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 13:54:05