Go语言sync.WaitGroup未生效 部分goroutine未执行完成就退出
Go sync.WaitGroup不生效问题排查与修复
根因分析
你的代码存在两个核心逻辑错误,直接导致WaitGroup未达预期效果:
- 通道收发次数不匹配:每个
IntradayStratifygoroutine处理单个标的时,会遍历所有时间周期(从代码看是5/15/30/60四个周期),每个周期无论是否匹配到策略都会往通道发送1条数据,但是你在主循环遍历标的时,每个标的仅从通道读取1次就处理下一个标的,剩余的发送操作直接阻塞在goroutine中,大量goroutine无法执行到defer wg.Done()逻辑。 - 主逻辑未等待任务执行完成:你单独启动了一个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
相关产品推荐
相关产品推荐

