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

使用Go的io.Copy拷贝网络流到文件时如何获取实时传输速度?

基于io.Reader包装的实时传输速率统计方案

直接自定义包装类实现io.Reader接口,在不改动原有io.Copy、文件写入、HTTP请求逻辑的前提下,就能拿到传输过程的实时字节速率,侵入性极低。

实现逻辑

  • 包装原始HTTP响应流,每次流读取操作时自动累加已传输的字节数
  • 启动独立协程按固定时间间隔采样,通过相邻两次采样的字节差除以间隔时长,得到实时字节/秒速率
  • 加读写锁规避多协程读写统计字段的数据竞争问题
  • 统计逻辑和业务逻辑完全解耦,回调函数支持自定义速率的后续处理(打印、上报、进度条更新都可以)

完整实现代码

package main

import (
	"fmt"
	"io"
	"net/http"
	"os"
	"sync"
	"time"
)

// RateReader 带速率统计的流包装器,实现io.Reader接口
type RateReader struct {
	reader         io.Reader    // 被包装的原始数据流
	totalBytes     int64        // 累计传输总字节数
	lastSnapshot   int64        // 上一次采样时的累计字节数
	lastSampTime   time.Time    // 上一次采样的时间
	currentRateBps int64        // 当前实时速率,单位:字节/秒
	mu             sync.RWMutex // 读写锁,规避数据竞争
}

// NewRateReader 初始化速率统计包装器
func NewRateReader(r io.Reader) *RateReader {
	return &RateReader{
		reader:       r,
		lastSampTime: time.Now(),
	}
}

// Read 实现io.Reader的读取方法,读取时同步累计字节
func (rr *RateReader) Read(p []byte) (n int, err error) {
	n, err = rr.reader.Read(p)
	if n > 0 {
		rr.mu.Lock()
		rr.totalBytes += int64(n)
		rr.mu.Unlock()
	}
	return
}

// StartCalc 启动速率定时统计
// interval: 采样间隔,建议设为1秒
// callback: 每次采样完成后的回调,会传入当前实时速率、累计传输总字节
func (rr *RateReader) StartCalc(interval time.Duration, callback func(rateBps, totalBytes int64)) {
	go func() {
		ticker := time.NewTicker(interval)
		defer ticker.Stop()
		for range ticker.C {
			rr.mu.RLock()
			currTotal := rr.totalBytes
			rr.mu.RUnlock()

			// 计算间隔内的传输速率
			byteDiff := currTotal - rr.lastSnapshot
			elapsed := time.Since(rr.lastSampTime).Seconds()
			rate := int64(float64(byteDiff) / elapsed)

			// 更新采样快照
			rr.mu.Lock()
			rr.currentRateBps = rate
			rr.lastSnapshot = currTotal
			rr.lastSampTime = time.Now()
			rr.mu.Unlock()

			// 触发回调
			if callback != nil {
				callback(rate, currTotal)
			}
		}
	}()
}

func main() {
	// 原有逻辑基本不需要改动
	res, err := http.Get("你的目标下载URL")
	if err != nil {
		panic(err)
	}

	out, err := os.OpenFile("output", os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644)
	if err != nil {
		panic(err)
	}

	defer out.Close()
	defer res.Body.Close()

	// 新增:包装响应体,启动速率统计
	rateR := NewRateReader(res.Body)
	rateR.StartCalc(time.Second, func(rateBps, totalBytes int64) {
		// 这里可以按需做单位转换,比如除以1024得KB/s,除以1024^2得MB/s
		fmt.Printf("\r实时速率: %d B/s | 已传输: %d 字节", rateBps, totalBytes)
	})

	// 原有io.Copy逻辑仅替换源为包装后的rateR即可
	_, err = io.Copy(out, rateR)
	if err != nil {
		panic(err)
	}
	fmt.Println("\n传输完成")
}

使用说明

  • 采样间隔可按需调整:需要更高实时性就缩短间隔,需要更平滑的速率曲线就把间隔设为1-2秒,避免过短间隔导致速率跳变太大
  • 如果需要计算全程平均速率,直接用总传输字节数除以总传输时长即可,不需要依赖定时采样逻辑
  • 这个包装器兼容所有实现io.Reader接口的流,本地文件拷贝、TCP流传输等场景都可以直接复用,不需要额外适配
  • 如果需要更平滑的滑动窗口速率,可以在结构体里加一个固定长度的环形队列,存储最近N次采样的字节数,取窗口内的平均值即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 19:54:17