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

Go语言多并发长运行任务进度追踪的正确实现方法

Go并行长任务进度追踪的惯用法实现

原实现的核心问题

你当前的写法不符合Go语言并发设计的惯用法,存在几个明确的缺陷:

  • 轮询逻辑存在天生阻塞:外层循环遍历所有tracker逐个select读通道的模式,只要轮询到暂时没有消息的通道,整个循环就会卡在该通道的读取上,这也是你给单个任务加延迟后出现全局阻塞的根本原因
  • 多通道拆分增加复杂度:为每个任务单独创建进度、错误、完成三个独立通道,很容易出现读写顺序错配、漏读消息、goroutine泄漏的问题
  • WaitGroup用法完全无效:代码里只调用了wg.Add,没有任何地方执行wg.Done(),单独起goroutine等待wg.Wait()的逻辑属于无效代码

符合Go惯用法的实现方案

核心思路非常简单:用统一事件通道替代分散的多通道,单消费端按事件到达顺序处理,完全避免轮询阻塞,具体实现逻辑如下:

  1. 定义统一的事件结构体,把进度更新、错误、任务完成三类状态全部封装在同一个结构里,附带任务唯一标识(比如你的场景里的URL)
  2. 所有并行worker共用一个全局事件通道,不需要为每个任务单独创建多个通道
  3. 正确使用WaitGroup管控worker生命周期:每个worker退出前必须调用wg.Done(),单独起一个goroutine等待所有worker执行完毕后关闭事件通道,作为消费端的退出信号
  4. 消费端直接用for range遍历事件通道,收到什么事件就处理什么,所有事件按到达顺序响应,完全不会被慢任务阻塞
  5. 全局总进度可以在消费端单goroutine内维护,不需要加锁,天然线程安全

修正后的完整实现代码

package main

import (
	"errors"
	"fmt"
	"strings"
	"sync"
	"time"
)

// Event 统一封装所有任务状态事件
type Event struct {
	Url      string // 任务标识
	Progress int    // 任务进度,取值0-100
	Error    error  // 任务执行错误,无错则为nil
	Done     bool   // 任务是否结束(包含成功、失败两种场景)
}

func work(url string, eventCh chan<- Event) {
	fmt.Printf("processing url %s\n", url)
	// 封装发送逻辑,避免重复写Url字段
	send := func(progress int, err error, done bool) {
		eventCh <- Event{
			Url:      url,
			Progress: progress,
			Error:    err,
			Done:     done,
		}
	}

	for i := 1; i <= 5; i++ {
		// 模拟单任务延迟场景
		if url == "google.com" {
			time.Sleep(time.Second * 3)
		}
		time.Sleep(time.Second)

		// 模拟.net站点出错场景
		if i == 3 && strings.HasSuffix(url, ".net") {
			send(0, errors.New("emulating error for .net sites"), true)
			return
		}

		send(20*i, nil, false)
	}

	// 任务正常完成
	send(100, nil, true)
}

func main() {
	urls := []string{"google.com", "youtube.com", "someurl.net"}
	eventCh := make(chan Event)
	var wg sync.WaitGroup

	// 启动所有并行任务
	for _, url := range urls {
		wg.Add(1)
		go func(u string) {
			defer wg.Done() // 无论任务成功失败,退出时必须标记完成
			work(u, eventCh)
		}(url)
	}

	// 等待所有任务结束后关闭事件通道,触发消费端退出
	go func() {
		wg.Wait()
		close(eventCh)
	}()

	// 消费端:单goroutine处理所有事件,无锁、无阻塞
	taskProgress := make(map[string]int, len(urls))
	for event := range eventCh {
		if event.Error != nil {
			fmt.Printf("Url = %s, error = %s\n", event.Url, event.Error.Error())
		}
		if event.Progress > 0 {
			taskProgress[event.Url] = event.Progress
			// 计算全局总进度
			total := 0
			for _, p := range taskProgress {
				total += p
			}
			fmt.Printf("Url = %s, progress = %d%%, total progress = %d%%\n",
				event.Url, event.Progress, total/len(urls))
		}
		if event.Done {
			fmt.Printf("Url = %s is completed\n", event.Url)
		}
	}

	fmt.Println("Everything is completed, exit")
}

方案优势

  • 彻底解决阻塞问题:所有事件按到达顺序即时处理,不会因为单个任务慢卡住所有状态更新
  • 架构简单易维护:不需要为每个任务维护多个通道,新增事件类型时只需要给Event结构体加字段即可,不需要改动核心调度逻辑
  • 天然线程安全:所有任务状态的修改都在消费端单goroutine内完成,不需要额外加锁
  • 完全符合Go并发设计哲学:通过通信传递状态,而非通过共享内存加锁访问

你之前考虑的“将ProgressTracker作为通道传入任务函数”的思路本质就是这个方案,只是之前把三类事件拆成三个独立通道走了弯路,把三类状态合并到同一个事件结构体走单通道,就能完全覆盖需求。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 01:15:38