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

Go并发目录大小计算程序出现goroutine死锁,求排查原因

目录大小计算程序死锁问题排查与修复

问题背景

编写了一个通过sizer结构体的Size方法递归计算目录总大小的程序,运行时触发All goroutines are asleep - deadlock.错误。尝试将无缓冲通道改为缓冲通道后问题仍未解决,程序通过cur.Ls(ctx)获取当前目录的子目录和文件列表,用dirQueue作为任务队列分发目录遍历任务。

死锁原因分析

  1. 初始发送操作阻塞主goroutine:Size方法中dirQueue <- d是无缓冲通道的发送操作,此时还未启动worker goroutine接收数据,主goroutine会直接卡在这个发送步骤,后续启动worker的代码完全无法执行(对应栈追踪中goroutine7的chan send阻塞)。
  2. 错误通道的阻塞式读取:Size方法一开始就执行err := <-errorQueue,但只有发生错误时worker才会向该通道发送数据,正常完成时无数据发送,导致主goroutine一直阻塞等待错误,无法推进后续逻辑。
  3. Worker队列关闭逻辑错误:多个worker都会检查len(dirQueue) == 0并尝试关闭队列,不仅会引发重复关闭通道的panic,而且该判断逻辑不可靠——判断完成后可能有其他worker刚向队列发送了新任务。
  4. 任务完成判断缺失:没有机制跟踪所有目录任务是否处理完毕,worker无法正确退出,WaitGroup也无法正确等待所有worker结束。

修复方案

修改后的完整代码

package storage

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

const MAX_GOROUTINES = 5

// Result represents the Size function result
type Result struct {
	// Total Size of File objects
	Size  int64
	// Count is a count of File objects processed
	Count int64
}

type DirSizer interface {
	// Size calculate a size of given Dir, receive a ctx and the root Dir instance
	// will return Result or error if happened
	Size(ctx context.Context, d Dir) (Result, error)
}

// sizer implement the DirSizer interface
type sizer struct {
	// maxWorkersCount number of workers for asynchronous run
	maxWorkersCount int
}

// NewSizer returns new DirSizer instance
func NewSizer() DirSizer {
	return &sizer{MAX_GOROUTINES}
}

func Worker(ctx context.Context, errorQueue chan error, dirQueue chan Dir, globalRes *Result,
	mutex *sync.Mutex, wg *sync.WaitGroup, taskCount *sync.WaitGroup) {
	defer wg.Done()
	for {
		select {
		case <-ctx.Done():
			return
		case cur, ok := <-dirQueue:
			if !ok {
				return
			}
			// 处理完当前目录,任务计数器减一
			defer taskCount.Done()

			dirs, files, err := cur.Ls(ctx)
			if err != nil {
				select {
				case errorQueue <- err:
				case <-ctx.Done():
				}
				return
			}
			res := Result{}
			for _, f := range files {
				fileSize, err := f.Stat(ctx)
				if err != nil {
					select {
					case errorQueue <- err:
					case <-ctx.Done():
					}
					return
				}
				res.Size += fileSize
				res.Count++
			}
			mutex.Lock()
			globalRes.Size += res.Size
			globalRes.Count += res.Count
			mutex.Unlock()
			// 添加子目录任务,计数器加对应数量
			taskCount.Add(len(dirs))
			for _, d := range dirs {
				select {
				case dirQueue <- d:
				case <-ctx.Done():
					return
				}
			}
		}
	}
}

func (a *sizer) Size(ctx context.Context, d Dir) (Result, error) {
	wg := sync.WaitGroup{}
	taskCount := sync.WaitGroup{}
	mutex := sync.Mutex{}
	errorQueue := make(chan error, 1) // 缓冲通道仅保留第一个错误
	dirQueue := make(chan Dir, a.maxWorkersCount) // 缓冲通道避免发送阻塞
	res := Result{}

	// 初始化根目录任务计数
	taskCount.Add(1)
	// 先启动所有worker
	for i := 0; i < a.maxWorkersCount; i++ {
		wg.Add(1)
		go Worker(ctx, errorQueue, dirQueue, &res, &mutex, &wg, &taskCount)
	}
	// 发送初始根目录任务
	dirQueue <- d

	// 单独goroutine等待所有任务完成后关闭队列
	go func() {
		taskCount.Wait()
		close(dirQueue)
	}()

	// 等待所有worker执行完毕
	wg.Wait()
	close(errorQueue)

	// 获取第一个错误(无错误则为nil)
	var err error
	select {
	case err = <-errorQueue:
	default:
	}
	fmt.Print(res)
	return res, err
}

关键修改说明

  • 任务计数器taskCount:用WaitGroup跟踪待处理的目录数量,每添加一个目录任务就Add(1),处理完成则Done(),计数器归零时关闭dirQueue,通知所有worker退出。
  • 通道操作非阻塞化:所有通道发送操作都结合select和ctx.Done(),避免上下文取消时出现阻塞。
  • 错误通道优化:使用缓冲大小为1的通道仅保留第一个错误,避免多错误发送导致阻塞,在所有worker结束后统一读取错误。
  • 初始化顺序调整:先启动worker再发送初始目录任务,彻底解决初始发送阻塞主goroutine的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 20:42:30