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

Golang中实现任一goroutine失败/超时终止全部的问题

问题解答

需求说明

  • 父任务包含3个子任务
  • 父任务设置5秒超时
  • 若所有子任务在父任务5秒超时时间内完成,程序正常终止
  • 若任一子任务在5秒超时时间内出错,需终止所有子进程
  • 若任一子任务未在5秒超时时间内完成,程序返回错误

代码与测试情况

我编写了如下Golang代码,用于测试第5种场景:将任务b的time.Sleep设置为10秒(超过父任务超时时间)。

package stackoverflow

import (
    "context"
    "fmt"
    "log"
    "sync"
    "testing"
    "time"

    "github.com/pkg/errors"
)

func TestTask(t *testing.T) {

    err := parent(context.Background())
    if err != nil {
        log.Printf("main error: %v\n", err)
    }
    log.Println("main done")

}

func parent(ctx context.Context) error {
    ctx, cancel := context.WithTimeout(ctx, time.Second*5)
    defer cancel()
    wg := &sync.WaitGroup{}

    fromChild := make(chan error)
    var err error
    wg.Add(1)
    go func(wg *sync.WaitGroup) {
        defer wg.Done()
        select {
        case err = <-fromChild:
            fmt.Printf("Parent: child go routines return error : %s \n", err.Error())
            return
        case <-ctx.Done():
            fmt.Printf("Parent: context ended. : %s \n", ctx.Err())
            return
        }
    }(wg)

    jobs := []string{"a", "b", "c"}
    for _, job := range jobs {
        wg.Add(1)
        go processChild(ctx, job, wg, fromChild)
    }
    wg.Wait()

    fmt.Printf("done.\n")
    if err != nil {
        return err
    }
    return nil
}
func processChild(ctx context.Context, name string, wg *sync.WaitGroup, ch chan<- error) {
    defer wg.Done()
    log.Printf("processsing %s job\n", name)

    c := time.After(time.Second * 60)
    go func() {
        select {
        case <-c:
            ch <- errors.New("error occurred.")
            return
        case <-ctx.Done():
            return
        }
    }()

    if name == "a" {
        time.Sleep(time.Second * 2)
    }
    if name == "b" {
        time.Sleep(time.Second * 10) //  when parent context timeout happens during this line, stop this goroutine.
    }
    log.Printf("done %s job\n", name)
}

测试输出如下,程序耗时10秒才结束,未按预期在5秒超时后终止任务b的goroutine:

$ go test -v . -count=1
=== RUN   TestTask
2023/11/21 10:57:08 processsing c job
2023/11/21 10:57:08 done c job
2023/11/21 10:57:08 processsing b job
2023/11/21 10:57:08 processsing a job
2023/11/21 10:57:10 done a job
Parent: context ended. : context deadline exceeded
done.
2023/11/21 10:57:18 main done
--- PASS: TestTask (10.00s) // it tooks 10s
PASS

期望与疑问

我期望processChild实现:

  1. 感知父context的超时,终止所有子goroutine
  2. 若子goroutine出错,返回错误并终止所有子goroutine
    请问该需求是否可实现?

解答

这个需求完全可以实现,当前代码的问题在于:

  1. time.Sleep是无感知阻塞调用,不会监听context的取消信号,所以即使父context超时,任务b的Sleep(10s)仍会执行完毕才退出
  2. 父goroutine的逻辑设计冗余,额外启动的监听goroutine没有起到主动终止所有子任务的作用
  3. 子任务内的超时监听goroutine逻辑无效,且未正确处理context取消

修正后的实现方案

核心思路:

  • 子任务通过select监听context取消信号,替换无感知的time.Sleep
  • 任一子任务出错或超时触发时,立即调用cancel()终止所有子任务
  • 简化错误传递和WaitGroup的管理逻辑

修正后的代码:

package stackoverflow

import (
	"context"
	"fmt"
	"log"
	"sync"
	"testing"
	"time"

	"github.com/pkg/errors"
)

func TestTask(t *testing.T) {
	err := parent(context.Background())
	if err != nil {
		log.Printf("main error: %v\n", err)
	}
	log.Println("main done")
}

func parent(ctx context.Context) error {
	ctx, cancel := context.WithTimeout(ctx, time.Second*5)
	defer cancel() // 确保函数退出时取消context

	wg := &sync.WaitGroup{}
	errChan := make(chan error, 1) // 带缓冲通道避免goroutine阻塞

	// 启动所有子任务
	jobs := []string{"a", "b", "c"}
	for _, job := range jobs {
		wg.Add(1)
		go processChild(ctx, job, wg, errChan)
	}

	// 等待第一个错误或超时信号
	select {
	case err := <-errChan:
		cancel() // 触发所有子任务终止
		wg.Wait() // 等待所有子任务退出
		return err
	case <-ctx.Done():
		wg.Wait()
		return ctx.Err()
	}
}

func processChild(ctx context.Context, name string, wg *sync.WaitGroup, errChan chan<- error) {
	defer wg.Done()
	log.Printf("processing %s job\n", name)

	// 定义各任务的执行时长
	var sleepDur time.Duration
	switch name {
	case "a":
		sleepDur = time.Second * 2
	case "b":
		sleepDur = time.Second * 10
	default:
		sleepDur = time.Second * 1
	}

	// 监听context取消与任务完成
	select {
	case <-ctx.Done():
		log.Printf("cancelled %s job\n", name)
		return
	case <-time.After(sleepDur):
		// 模拟子任务出错场景
		if name == "a" {
			select {
			case errChan <- errors.New(fmt.Sprintf("%s job failed", name)):
			default: // 避免通道满时阻塞
			}
			return
		}
		log.Printf("done %s job\n", name)
	}
}

关键改进点

  1. 可中断任务逻辑:用select结合time.After和ctx.Done()替代time.Sleep,context取消时能立即终止子任务
  2. 全局终止触发:父goroutine收到错误或超时信号后,立即调用cancel()终止所有子任务
  3. 避免阻塞风险:使用带缓冲的errChan,并通过select+default确保错误发送不会阻塞子goroutine
  4. 简化流程:移除冗余的监听goroutine,父goroutine直接通过select等待终止信号,逻辑更清晰

测试该代码,超时场景下程序会在5秒左右结束,错误场景下会立即终止所有子任务,完全符合需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 15:04:58