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

Go语言REST请求缓存文件内容不完整问题排查

批量REST接口缓存文件随机不完整/开头缺失问题

问题描述

我编写了一个批量调用REST接口的Go程序,将请求结果缓存到本地文件后返回文件指针供调用方使用。99%的情况下代码运行正常,但会随机出现缓存文件未包含完整响应体的情况,有时甚至是文件开头缺失、结尾正常。调用方仅通过readFromCache()方法调用。

原始代码

file, err := CacheProxy{}.readFromCache(url, path.Join("CACHE", a.Hostname, a.ID, s))
func (c CacheProxy) readFromCache(url string, p string) (io.ReadCloser, error) {
    c.makeDirectory(p)
    file, err := os.Open(p)
    if err != nil {
        // 文件不存在时,从URL读取并缓存
        return c.readFromURL(url, p)
    }
    return file, nil
}

func (c CacheProxy) readFromURL(url string, p string) (io.ReadCloser, error) {
    client := http.Client{}
    resp, err := client.Get(url)
    if err != nil {
        log.Errorf("Cant read from URL : %s", err)
        return nil, err
    }

    defer resp.Body.Close()

    out, err := os.Create(p)
    if err != nil {
        log.Errorf("Unable to open file : %s", err)
        return nil, err
    }

    writer := bufio.NewWriter(out)
    _, err = io.Copy(writer, resp.Body)
    if err != nil {
        log.Errorf("Unable to copy file : %s", err)
        return nil, err
    }

    writer.Flush()

    _, err = out.Seek(0, 0)
    if err != nil {
        log.Errorf("Unable seek to beginning : %s", err)
        return nil, err
    }
    return out, nil
}

尝试过的方案

更新1:全量加载到内存后写入(无问题但占用内存高)

func (c CacheProxy) readFromURL(url string, p string) (io.ReadCloser, error) {
    resp, err := http.Get(url)
    if err != nil {
        log.Errorf("Cant read from URL : %s", err)
        return nil, err
    }

    defer resp.Body.Close()

    if resp.StatusCode != 200 {
        return nil, fmt.Errorf("ResultCode != 200 (%d), aborting", resp.StatusCode)
    }

    body, err := io.ReadAll(resp.Body)
    if err != nil {
        log.Errorf("Error while reading body : %s", err)
        return nil, err
    }

    os.WriteFile(p, body, 0600)

    out, err := os.Open(p)
    if err != nil {
        log.Errorf("Cannot reopen file : %s", err)
        return nil, err
    }
    return out, nil
}

更新2:改用file.ReadFrom()(仍存在问题)

func (c CacheProxy) readFromURL(url string, p string) (io.ReadCloser, error) {
    resp, err := http.Get(url)
    if err != nil {
        log.Errorf("Cant read from URL : %s", err)
        return nil, err
    }

    defer resp.Body.Close()

    if resp.StatusCode != 200 {
        return nil, fmt.Errorf("ResultCode != 200 (%d), aborting", resp.StatusCode)
    }

    out, err := os.Create(p)
    if err != nil {
        log.Errorf("Unable to open file : %s", err)
        return nil, err
    }
    out.ReadFrom(resp.Body)
    if err != nil {
        log.Errorf("Unable to read stream : %s", err)
        return nil, err
    }
    err = out.Close()
    if err != nil {
        log.Errorf("Error while closing file : %s", err)
        return nil, err
    }

    out, err = os.Open(p)
    if err != nil {
        log.Errorf("Cannot reopen file : %s", err)
        return nil, err
    }
    return out, nil
}

更新3:问题现象补充

有时缓存文件并非被截断,而是开头缺失、结尾正常。排除同一文件被重复写入覆盖的可能,http.Client已确认线程安全。


问题原因分析

核心问题是缓存检查与写入的竞态条件:
当多个线程同时处理同一个缓存路径p时,第一个线程进入readFromURL并通过os.Create创建了文件,但还未完成写入;此时第二个线程调用os.Open会成功(因为文件已存在),直接返回这个还在写入中的文件句柄,导致读取到不完整或开头缺失的内容。

另外,更新2中的代码存在明显错误:调用out.ReadFrom(resp.Body)时未捕获其返回的错误,后续的错误判断仍使用之前os.Create的错误值,无法检测到写入过程中出现的问题。

高效解决方案

方案1:基于缓存路径的互斥锁

使用sync.Map存储每个缓存路径的互斥锁,确保同一时间只有一个线程处理同一个缓存文件的写入:

var cacheLocks sync.Map

func (c CacheProxy) readFromCache(url string, p string) (io.ReadCloser, error) {
    // 获取或创建当前缓存路径的锁
    lockIface, _ := cacheLocks.LoadOrStore(p, &sync.Mutex{})
    lock := lockIface.(*sync.Mutex)
    lock.Lock()
    defer lock.Unlock()

    c.makeDirectory(p)
    file, err := os.Open(p)
    if err != nil {
        return c.readFromURL(url, p)
    }
    return file, nil
}

// 修正后的readFromURL:写入完成后关闭文件,重新打开返回
func (c CacheProxy) readFromURL(url string, p string) (io.ReadCloser, error) {
    client := http.Client{}
    resp, err := client.Get(url)
    if err != nil {
        log.Errorf("Cant read from URL : %s", err)
        return nil, err
    }
    defer resp.Body.Close()

    if resp.StatusCode != http.StatusOK {
        return nil, fmt.Errorf("ResultCode != 200 (%d), aborting", resp.StatusCode)
    }

    out, err := os.Create(p)
    if err != nil {
        log.Errorf("Unable to open file : %s", err)
        return nil, err
    }
    defer out.Close()

    _, err = io.Copy(out, resp.Body)
    if err != nil {
        log.Errorf("Unable to copy file : %s", err)
        // 删除不完整的缓存文件
        os.Remove(p)
        return nil, err
    }

    // 强制刷入磁盘
    if err = out.Sync(); err != nil {
        log.Errorf("Failed to sync file : %s", err)
        os.Remove(p)
        return nil, err
    }

    // 重新打开文件返回读句柄
    return os.Open(p)
}

方案2:临时文件+原子重命名

先将响应写入临时文件,完成后原子重命名到目标路径,避免其他线程读取到不完整的文件:

func (c CacheProxy) readFromCache(url string, p string) (io.ReadCloser, error) {
    c.makeDirectory(p)
    file, err := os.Open(p)
    if err != nil {
        return c.readFromURL(url, p)
    }
    return file, nil
}

func (c CacheProxy) readFromURL(url string, p string) (io.ReadCloser, error) {
    client := http.Client{}
    resp, err := client.Get(url)
    if err != nil {
        log.Errorf("Cant read from URL : %s", err)
        return nil, err
    }
    defer resp.Body.Close()

    if resp.StatusCode != http.StatusOK {
        return nil, fmt.Errorf("ResultCode != 200 (%d), aborting", resp.StatusCode)
    }

    // 创建临时文件,使用p+固定tmp后缀
    tmpPath := p + ".tmp"
    out, err := os.OpenFile(tmpPath, os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0600)
    if err != nil {
        // 如果临时文件已存在,说明有其他线程在处理,短暂等待后重试
        log.Errorf("Unable to create temp file : %s", err)
        time.Sleep(100 * time.Millisecond)
        return c.readFromCache(url, p)
    }
    defer func() {
        // 无论成功失败,都清理临时文件
        os.Remove(tmpPath)
    }()

    _, err = io.Copy(out, resp.Body)
    if err != nil {
        log.Errorf("Unable to copy file : %s", err)
        out.Close()
        return nil, err
    }

    if err = out.Sync(); err != nil {
        log.Errorf("Failed to sync file : %s", err)
        out.Close()
        return nil, err
    }

    if err = out.Close(); err != nil {
        log.Errorf("Failed to close temp file : %s", err)
        return nil, err
    }

    // 原子重命名,确保其他线程看到的是完整的文件
    if err = os.Rename(tmpPath, p); err != nil {
        log.Errorf("Failed to rename temp file : %s", err)
        return nil, err
    }

    return os.Open(p)
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 05:08:06