Kratos 云盘项目:对象存储与文件分片上传

2025-01-15T10:30:00+08:00 | 36分钟阅读 | 更新于 2025-01-15T10:30:00+08:00

@

学习目标

学完本文你应该能够:

  1. 说清存储抽象的价值:用接口隔离 biz 层与具体存储实现(本地磁盘 / MinIO / OSS),并能在面试里讲出"编译期接口校验"这一加分点。
  2. 讲透秒传与分片上传:从 FindByHash 秒传,到 InitUpload → UploadPart → MergeParts 的完整链路,包括分片大小钳制、分片校验。
  3. 解释断点续传与合并的边界条件FindByHashAndSize 如何恢复会话、CheckParts 如何让客户端补传、MergeParts 如何排序合并并 Sync 落盘。
  4. 理解一致性与资源清理:物理文件与 DB 记录的回滚、孤儿分片清理、上传并发信号量限流。
  5. 掌握流式下载与 MIME 探测io.ReadCloser 零拷贝、GetFileURL 的 presigned URL、内置 builtinMimeTypes 兜底映射与危险扩展名拦截。

前置知识:Go 的 io.Reader / io.ReadCloser、接口与"鸭子类型"、Kratos 分层(biz / data / service)、HTTP 文件上传(multipart)、基础文件系统(os / bufio)。

动手 3 件

  • 把这个项目的 Storage 接口抄一遍,自己写个 memoryStorage 实现并加上 var _ biz.Storage = (*memoryStorage)(nil) 验证编译期实现。
  • curl 模拟一次分片上传:先 InitUpload,再分两次 UploadPart,最后 MergeParts,观察磁盘上 chunks/<uploadID>/part_00001 的生成与合并。
  • 故意让 MergeParts 成功但 DB 落库失败,确认代码会删除已合并的物理文件(回滚),不产生孤儿文件。

一、为什么需要存储抽象:Storage 接口

Q1. 什么是对象存储?为什么微服务里要把"存文件"抽象成一个接口?

答:

类比:你家里的电器(电脑、手机、台灯)都依赖"墙上的电源插座"。插座只规定了"两孔/三孔 + 220V"这个契约,电器并不关心发电站是火电、水电还是核电。哪天你换了供电公司,电器照常工作。对象存储也是一样——业务层(云盘的"上传文件"用例)只关心"我能把一个 io.Reader 存进去,并能用一个 storagePath 把它读出来",至于底层是本地磁盘、MinIO 还是阿里云 OSS,业务层根本不该知道。

工程内容:本项目在 internal/biz/storage.go 里定义了 Storage 接口,它是整个文件子系统的"插座标准"。biz.FileUsecase 只依赖这个接口,不 import 任何具体存储包:

// internal/biz/storage.go
type Storage interface {
    Upload(ctx context.Context, filePath string, reader io.Reader) (storagePath string, err error)
    Download(ctx context.Context, storagePath string) (io.ReadCloser, error)
    Delete(ctx context.Context, storagePath string) error
    InitMultipartUpload(ctx context.Context, uploadID, fileName string) error
    UploadPart(ctx context.Context, uploadID string, partNumber int, reader io.Reader) (PartInfo, error)
    CompleteMultipartUpload(ctx context.Context, uploadID string, parts []PartInfo) (storagePath string, err error)
    AbortMultipartUpload(ctx context.Context, uploadID string) error
    CleanupOrphanChunks(ctx context.Context, maxAge time.Duration) ([]string, error)
    GetFileURL(ctx context.Context, storagePath string, expire time.Duration) (string, error)
}

注意 PartInfo 这个元数据载体:

// internal/biz/storage.go
type PartInfo struct {
    PartNumber int
    Hash       string
    Size       int64
}

接口把"存储"这件事收敛成了 9 个方法,覆盖:小文件上传(Upload)、下载(Download)、删除(Delete)、分片上传的初始化/逐片/合并/中止/清孤儿(Init/UploadPart/Complete/Abort/CleanupOrphanChunks)、以及取访问 URL(GetFileURL)。

另外值得强调的是,接口方法返回的是 storagePath(一个能唯一定位对象的字符串,本地是绝对路径、MinIO 是 年/月/日/uuid.ext 对象键),而不是 URLio.Reader。这个设计很关键:存储路径是稳定的内部标识,访问 URL 是临时的、可派生的GetFileURL 单独存在,意味着"文件在哪"和"文件怎么给人看"被拆开——本地存储把路径包成 file://,MinIO/OSS 则按需签发限时签名 URL。这样换存储实现时,DB 里存的 Path 语义不变,只有"如何生成可访问链接"随实现而变。

回到 Storage 接口本身,九个方法可以分成四组理解:① 基础 CRUD——Upload/Download/Delete;② 分片生命周期——InitMultipartUpload/UploadPart/CompleteMultipartUpload/AbortMultipartUpload;③ 运维兜底——CleanupOrphanChunks;④ 访问派生——GetFileURL。biz 层只编排这些方法,从不过问"底层到底是磁盘还是对象存储",这正是依赖倒置原则(DIP)在 Go 接口上的落地。顺带一提,biz 层不依赖具体实现的另一个好处是单测可以轻松注入一个内存版 storage 实现,把上传与合并逻辑测个遍而不碰真实磁盘。

下面这张图展示 biz 层如何通过接口与三种实现解耦:

flowchart TD
    A["biz.FileUsecase
上传/下载/秒传/分片"] -->|依赖| B["biz.Storage 接口
9 个方法"] B -->|实现| C["localStorage
本地磁盘 ./data/uploads"] B -->|实现| D["minioStorage
S3 兼容 对象存储"] B -->|实现| E["ossStorage
阿里云 OSS stub"] F["NewLocalStorage / NewMinioStorage
/ NewOssStorage"] -->|注入| B

三种实现在 internal/data/storage/ 下分别放在 local.gominio.gooss_stub.go。切换实现只是 wiremain 里换一个 NewXxxStorage 工厂,业务代码一行都不用改。

Q2. 怎么保证"实现类真的实现了接口"?OSS 为什么只写了 stub?

答:

类比:买插头前你会先在墙上的插座上试一下,确认孔位对得上。Go 里不能在编译期"试",但有一个惯用法能做等价检查——用一个不会被使用的变量声明来"骗"编译器做类型断言。

工程内容:三个实现文件开头都有这样一行:

// internal/data/storage/local.go
var _ biz.Storage = (*localStorage)(nil)

// internal/data/storage/minio.go
var _ biz.Storage = (*minioStorage)(nil)

// internal/data/storage/oss_stub.go
var _ biz.Storage = (*ossStorage)(nil)

var _ biz.Storage = (*localStorage)(nil) 的意思是:把 *localStorage 的零值赋给一个 biz.Storage 类型的匿名变量。如果 localStorage 少了某个方法,这行直接编译失败。这是面试里非常加分的细节——它把"接口是否被实现"的检查从运行时提前到了编译期,而且不影响运行时(变量名为 _,会被优化掉)。

至于 OSS,项目里用 oss_stub.go 故意只留壳:

// internal/data/storage/oss_stub.go
type ossStorage struct {
    // 接入后填:client *oss.Client, bucket *oss.Bucket
}

func (s *ossStorage) Upload(ctx context.Context, filePath string, reader io.Reader) (string, error) {
    // TODO: bucket.PutObject(filePath, reader)
    return "", biz.ErrStorageNotImplemented
}
// ... 其余方法全部返回 biz.ErrStorageNotImplemented

每个方法都给了 // TODO 注释指明接入点(例如 GetFileURL 里写 bucket.SignURL(...))。这样做的好处:本地开发与单测可以先跑 MinIO / 本地磁盘,等真正要上阿里云时再补实现,不会因为"还没接 OSS"就挡住整体编译和交付。注释里还标了签名 URL、分片上传等对接点,新人接手时能照着填空。

⚠️ :如果给 OSS stub 的实现留空却不返回 ErrStorageNotImplemented,而是返回 nil,那么调用方会以为上传成功了,实际磁盘/对象存储里什么都没有,最后 DB 里会多出一条指向空 Path 的"幽灵文件"。stub 一定要显式返回"未实现"错误,而不是默默成功。


二、小文件上传与按日期分目录

Q3. 小文件是怎么落盘的?为什么要用"按日期分目录"?

答:

类比:把一年的邮件全塞进同一个抽屉,找一封要翻到地老天荒,而且抽屉一旦损坏全部完蛋。聪明做法是一月一个文件夹、一天一个子文件夹。文件系统也是同理——一个目录里塞几百万个文件,很多文件系统(尤其是 ext4 默认配置、ReiserFS 之外的实现)的 readdir / 创建性能会显著下降。

工程内容localStorage.UploaddatePath() 把文件分散到 uploadDir/年/月/日/ 下:

// internal/data/storage/local.go
func (s *localStorage) datePath() string {
    now := time.Now()
    return filepath.Join(s.uploadDir,
        fmt.Sprintf("%d", now.Year()),
        fmt.Sprintf("%02d", now.Month()),
        fmt.Sprintf("%02d", now.Day()),
    )
}

func (s *localStorage) Upload(ctx context.Context, filePath string, reader io.Reader) (string, error) {
    dir := s.datePath()
    if err := os.MkdirAll(dir, 0755); err != nil {
        return "", fmt.Errorf("storage: failed to create upload dir: %w", err)
    }
    destPath := filepath.Join(dir, filePath) // filePath 形如 "uuid.extension"
    f, err := os.Create(destPath)
    if err != nil {
        return "", fmt.Errorf("storage: failed to create file: %w", err)
    }
    defer f.Close()
    if _, err := io.Copy(f, reader); err != nil { // 流式写入,不整文件加载内存
        return "", fmt.Errorf("storage: failed to write file: %w", err)
    }
    return destPath, nil
}

两个关键点:

  1. io.Copy(f, reader) 流式落盘:数据从 HTTP 请求体(一个 io.Reader)直接写到磁盘文件,全程不把整份文件读进内存。这对大文件极其重要——否则一个 2GB 文件上传就会吃掉 2GB 内存。
  2. filePath 由调用方传入 uuid + ext:文件名用 newUUID() 生成(internal/biz/file.go),避免用户传来的中文名/特殊字符污染磁盘,也避免同名覆盖。

关于"为什么不把整文件读进内存",这里展开一下:io.Copy 的底层实现是从 src 读取一块(默认 32KB)写到 dst,循环直到 EOF。它对源和目的都只要求实现 io.Reader / io.Writer完全不感知"整份文件有多大"。这意味着:无论上传的是 1KB 文本还是 5GB 视频,单请求常驻内存都只是那块 buffer,而不是整文件。对比之下,如果业务里写成 data, _ := io.ReadAll(r); storage.Upload(...),内存峰值就是文件大小——并发 100 个 2GB 上传,OOM 是必然。所以一句话总结:凡是"文件流转"场景,永远用 io.Reader/io.Writer 管道,不要用 []byte 全量缓冲。本项目 bytesToReader 只是小文件便捷封装,真正的大文件走的是 HTTP multipart 的 Reader 直传。

调用方在 FileUsecase.Upload 里先算好 fileName

// internal/biz/file.go(小文件上传用例,节选)
uuid, _ := newUUID()
ext := extractExt(name)
fileName := uuid + ext
storagePath, err := uc.storage.Upload(ctx, fileName, bytesToReader(data))
if err != nil {
    return nil, err
}
file := &File{ /* ... Path: storagePath ... */ }
created, err := uc.fileRepo.Create(ctx, file)
if err != nil {
    // 数据库写入失败,回滚已上传的存储文件,避免孤儿文件
    if delErr := uc.storage.Delete(ctx, storagePath); delErr != nil {
        log.Error("file: failed to rollback storage upload after db create failed", ...)
    }
    return nil, err
}

物理文件与 DB 一致性是面试高频:

  • 上传成功、写 DB 失败 → 回滚删物理文件(如上面代码,避免磁盘上留下没人引用的孤儿文件)。
  • 删除文件的顺序 → 先删物理成功,再标记 DB:因为物理文件删了还能从 DB 恢复引用,但 DB 先删了物理还在,就成了永久孤儿。本文后面"清理"小节会再强调。

下面这张图把小文件上传与"DB 失败回滚"画出来:

flowchart TD
    A["客户端 小文件字节"] --> B["FileUsecase.Upload
检查配额 + 生成 uuid"] B --> C["storage.Upload
io.Copy 流式写盘"] C -->|成功| D["fileRepo.Create 写 DB"] D -->|失败| E["storage.Delete
删除物理文件回滚"] D -->|成功| F["AddUsedStorage 加配额
写入文件列表缓存"] E --> G["返回错误 无孤儿文件"] F --> H["返回 File 记录"]

⚠️ io.Copy 之后忘了 f.Sync()?小文件上传通常不强制 Sync,因为 OS 页缓存会在后台刷盘;但分片合并(下面讲)必须 Sync,否则进程崩溃可能丢失尚未落盘的数据。两者的可靠性要求不同,不要混为一谈。


三、秒传:同样的文件只存一份

Q4. 秒传(instant upload)是怎么做到的?

答:

类比:公司群里有人发了份 2GB 的安装包,你也想发同样一份。聪明做法是:你先把文件的"指纹"(SHA-256)发出来问一句"谁已经有了?",群友说"我有了",你就不传了,直接在共享盘里"建一条指向同一份文件的快捷方式"。秒传就是这个意思——传指纹不传内容

工程内容:客户端在上传前对本地文件算 SHA-256,把 hash 带给服务端。FileUsecase.SecUploadFindByHash 查同用户下是否已有相同哈希的文件:

// internal/biz/file.go
func (uc *FileUsecase) SecUpload(ctx context.Context, userID uint64, hash, name string, parentID *uint64) (*File, bool, error) {
    name = sanitizeFileName(name)
    if isDangerousFileType(name) {
        return nil, false, perrors.BadRequest("FILE_TYPE_FORBIDDEN", "不允许上传此类型的文件")
    }

    existing, err := uc.fileRepo.FindByHash(ctx, userID, hash)
    if err != nil {
        if perrors.IsNotFound(err) {
            return nil, false, nil // 没命中,让调用方走真实上传
        }
        return nil, false, err
    }

    // 继续前检查存储空间
    if uc.userUC != nil {
        if err := uc.userUC.CheckStorageAvailable(ctx, userID, existing.Size); err != nil {
            return nil, false, err
        }
    }

    uuid, _ := newUUID()
    file := &File{
        UUID:     uuid,
        UserID:   userID,
        Name:     name,
        Size:     existing.Size,
        Type:     existing.Type,
        Hash:     existing.Hash,
        Path:     existing.Path, // 关键:指向同一个物理文件
        ParentID: parentID,
        Status:   0,
    }
    created, err := uc.fileRepo.Create(ctx, file)
    if err != nil {
        return nil, false, err
    }
    // 成功后增加已用存储
    if uc.userUC != nil {
        if err := uc.userUC.AddUsedStorage(ctx, userID, existing.Size); err != nil {
            log.Error(...)
        }
    }
    return created, true, nil
}

要点:

  1. 命中后新建的 File 记录 Path 直接复用 existing.Path,物理文件一点没动,零拷贝、零重复存储。
  2. 配额仍要 +Size:虽然没多存物理字节,但"这个用户的可用空间"要按逻辑占用记账,否则用户能无限"秒传"撑爆配额统计。
  3. 危险类型先拦截isDangerousFileType(name) 在查 hash 之前就挡掉 .exe / .sh / .php 等(见第七节)。

此外还有一个工程细节值得在面试里主动提:秒传让多个 File 记录指向同一份物理文件,因此"删除"必须做引用计数。本项目 SecUpload 只新增记录、不复制物理字节,那么当其中一个用户删除文件时,不能无脑 storage.Delete(Path)——否则别人的秒传记录会指向一个已消失的文件。正确做法是:删除前先查"还有多少 File 记录引用同一 Path",引用数降到 0 才真正删物理文件(本项目删除链路虽不在本次必读范围,但这是秒传设计的必然后果,面试常被追问)。这也解释了为什么配额要按"逻辑占用"记账而非"物理字节去重"——两个用户各秒传了一份 2GB 电影,物理上只占 2GB,但各自配额都应扣 2GB,否则 A 删了 B 的配额凭空多出来。

⚠️ :秒传依赖"SHA-256 相同即文件相同"的假设。它正确,但要注意哈希碰撞攻击恶意预占位:攻击者可以先用小体积构造一个特定 hash 的文件"秒传占位",诱导别人以为传了同一个文件。工程上通常再结合 Size 一起校验(本项目 FindByHashAndSize 就是这么做的,见断点续传),并对单物理文件做引用计数,删除时"引用归零才真删"。

秒传的完整判定流程如下:

flowchart TD
    A["客户端 计算 SHA-256"] --> B["SecUpload 带 hash 进来"]
    B --> C["isDangerousFileType 拦截?"]
    C -->|危险| D["拒绝 400"]
    C -->|安全| E["fileRepo.FindByHash"]
    E -->|未命中| F["返回 false
走真实上传/分片"] E -->|命中| G["CheckStorageAvailable 配额"] G -->|不足| H["413 STORAGE_INSUFFICIENT"] G -->|够| I["新建 File 记录
Path 复用 existing.Path"] I --> J["AddUsedStorage 记账"] J --> K["返回 true 秒传成功"]

四、分片上传:大文件的稳妥之道

Q5. 分片上传的完整流程是什么?合并时为什么要 1MB 缓冲 + Sync?

答:

类比:寄一台大冰箱,快递公司不让整车发,让你拆成 12 个箱子分别寄,每个箱子编号 1~12。收件人收到后按编号顺序拼回去,拼完确认无误,再把 12 个箱子销毁。分片上传一模一样——把大文件切成 N 片,每片独立上传、独立重试,最后服务端按编号顺序合并。

工程内容:链路是 InitUpload → UploadPart* → MergeParts(失败用 AbortMultipartUpload)。

1) InitUpload(biz 层):先钳制分片大小,生成 uploadID用 biz 自己生成的 uploadID 去初始化 storage 的分片目录,这是避免孤儿目录的关键:

// internal/biz/file.go(节选)
chunksTotal := int32((fileSize + chunkSize - 1) / chunkSize)
uploadID := fmt.Sprintf("upload_%s", generateToken())

// 使用 biz 层生成的 uploadID 初始化 storage 层分片目录,
// 确保 biz 会话 ID 与 storage 目录名一致,避免孤儿空目录
if err := uc.storage.InitMultipartUpload(ctx, uploadID, fileName); err != nil {
    return nil, err
}
session := &UploadSession{
    UploadID: uploadID, /* ... */ Status: 0,
}
created, err := uc.uploadRepo.CreateSession(ctx, session)
if err != nil {
    uc.storage.AbortMultipartUpload(ctx, uploadID) // DB 失败则回滚清理目录
    return nil, err
}

localStorage.InitMultipartUpload 建的就是 chunks/<uploadID> 目录:

// internal/data/storage/local.go
func (s *localStorage) InitMultipartUpload(ctx context.Context, uploadID, fileName string) error {
    chunkDir := filepath.Join(s.uploadDir, "chunks", uploadID)
    if err := os.MkdirAll(chunkDir, 0755); err != nil {
        return fmt.Errorf("storage: failed to create chunk dir: %w", err)
    }
    return nil
}

2) UploadPart:每片写成 part_%05d(5 位零填充编号),并先 Stat 防止重复上传同一片:

// internal/data/storage/local.go
func (s *localStorage) UploadPart(ctx context.Context, uploadID string, partNumber int, reader io.Reader) (biz.PartInfo, error) {
    chunkDir := filepath.Join(s.uploadDir, "chunks", uploadID)
    if err := os.MkdirAll(chunkDir, 0755); err != nil {
        return biz.PartInfo{}, fmt.Errorf(...)
    }
    partPath := filepath.Join(chunkDir, fmt.Sprintf("part_%05d", partNumber))
    if _, err := os.Stat(partPath); err == nil {
        return biz.PartInfo{}, biz.ErrPartAlreadyUploaded // 已存在直接报错,支持幂等补传
    }
    f, _ := os.Create(partPath)
    defer f.Close()
    written, err := io.Copy(f, reader)
    // ...
    return biz.PartInfo{PartNumber: partNumber, Size: written}, nil
}

part_%05d 这个命名很关键:零填充保证字典序 == 数字序,后续按文件名 ls 也能天然有序,且合并时 sort 不依赖天然顺序(仍显式排序,双保险)。

3) MergeParts:先按 PartNumber 排序,再用 bufio.NewWriterSize(dest, 1024*1024) 1MB 缓冲顺序合并,最后 Flush + dest.Sync() 落盘,再 os.RemoveAll(chunkDir) 清理分片目录:

// internal/data/storage/local.go(节选)
sort.Slice(parts, func(i, j int) bool {
    return parts[i].PartNumber < parts[j].PartNumber
})
dest, _ := os.Create(finalPath)
defer dest.Close()

bufWriter := bufio.NewWriterSize(dest, 1024*1024) // 1MB 缓冲,减少 write 系统调用
defer bufWriter.Flush()
copyBuf := make([]byte, 1024*1024) // 复用缓冲区,避免每次 io.Copy 重新分配

for _, part := range parts {
    partPath := filepath.Join(chunkDir, fmt.Sprintf("part_%05d", part.PartNumber))
    src, _ := os.Open(partPath)
    if _, err := io.CopyBuffer(bufWriter, src, copyBuf); err != nil {
        src.Close()
        return "", fmt.Errorf("storage: failed to merge part %d: %w", part.PartNumber, err)
    }
    src.Close()
}
if err := bufWriter.Flush(); err != nil { /* ... */ }
if err := dest.Sync(); err != nil {           // 强制刷盘,防进程崩溃丢数据
    return "", fmt.Errorf("storage: failed to sync merged file: %w", err)
}
os.RemoveAll(chunkDir) // 合并成功后删除分片目录
return finalPath, nil

为什么用 1MB 缓冲 + Sync

  • bufio.NewWriterSize 把"每次写一块就触发一次 write 系统调用"聚合成"攒满 1MB 才刷一次",对几百个分片的大文件,系统调用次数从"分片数 × 分片内次数"降到约"最终文件大小 / 1MB",吞吐显著提升。
  • dest.Sync() 把内核页缓存里的数据真正写到磁盘介质,防止"程序刚返回成功、机器就断电"导致文件残缺。合并是一次性不可逆操作,必须确保落盘后再告诉客户端"上传完成"。
  • copyBuf 复用同一块 1MB 内存做 io.CopyBuffer 的搬运缓冲,避免 Go 默认 32KB 缓冲在超大文件下反复分配,也避免用超大缓冲占用内存。

顺带对比 MinIO 实现(minio.go),能看出"同一接口、不同实现"的威力。minioStorage 的分片合并不自己读文件,而是调用 client.ComposeObject——由对象存储服务端把 chunks/{uploadID}/part_xxxxx 这些对象在服务端直接拼成最终对象,省去了"下载到本地再上传"的双重流量。InitMultipartUpload 在 MinIO 下只是放一个 chunks/{uploadID}/.init 标记对象来证明会话存在;合并成功后清理分片用的是 go s.cleanupChunks(context.Background(), uploadID)——特意用 context.Background() 而非请求 ctx,避免客户端断开导致清理中断留下孤儿。这与本地存储"同步 RemoveAll“形成鲜明对比:本地合并完当场删目录,MinIO 则异步、带独立 context 删对象。两种方式都正确,差别只在于"本地文件系统删除是毫秒级、可同步;对象存储批量删除可能稍慢、宜异步”。

合并流程的时序如下(请求→biz→storage→磁盘):

flowchart TD
    A["客户端 MergeParts"] --> B["biz.MergeParts
校验分片数 == ChunksTotal"] B --> C["storage.CompleteMultipartUpload"] C --> D["sort 按 PartNumber 排序"] D --> E["bufio 1MB 缓冲
顺序 io.CopyBuffer 合并"] E --> F["bufWriter.Flush"] F --> G["dest.Sync 强制落盘"] G --> H["os.RemoveAll chunks 目录"] H --> I["返回 finalPath"] I --> J["biz 写 File 记录 + AddUsedStorage"] J --> K["发布 UploadCompleted 事件"]

Q6. 分片大小为什么要钳制在 [1MB, 100MB],文件上限 10GB?

答:

类比:快递公司规定"每个箱子至少 1kg、最多 100kg,单票总重不超过 10 吨"。如果允许 1 克的箱子,一趟寄 100 万箱,分拣员累死(分片数爆炸);如果允许 1 吨的箱子,单个箱子就撑爆货车(内存/磁盘瞬时压力);如果不限总重,仓库被一个人占满(资源被独占)。

工程内容internal/biz/file.go 顶部有四个常量,是服务端的"硬杠杠":

// internal/biz/file.go
const uploadSemaphoreSize = 50

const (
    minChunkSize = 1 * 1024 * 1024         // 1MB  分片下限
    maxChunkSize = 100 * 1024 * 1024       // 100MB 分片上限
    maxFileSize  = 10 * 1024 * 1024 * 1024 // 10GB 单文件硬上限
)

InitUpload 里对客户端传来的 chunkSize 做钳制:

// internal/biz/file.go
if fileSize > maxFileSize {
    return nil, perrors.BadRequest("FILE_TOO_LARGE", "文件大小超过系统上限")
}
if chunkSize <= 0 {
    chunkSize = 5 * 1024 * 1024 // 默认 5MB
}
if chunkSize < minChunkSize {
    chunkSize = minChunkSize
}
if chunkSize > maxChunkSize {
    chunkSize = maxChunkSize
}

为什么必须钳制(而不是信任客户端):

  • chunkSize=1 制造天文数字分片:一个 10GB 文件按 1 字节分片 = 100 亿片,每片在磁盘/对象存储里都是一次写入 + 一条 DB 记录,ChunksTotal 直接溢出、会话表被塞爆。下限 1MB 把最大分片数压到约 1 万片(10GB / 1MB)。
  • 防超大分片撑爆内存:如果客户端说"我这一片 2GB",服务端在合并/校验时可能尝试一次性分配或缓存该分片。上限 100MB 把单分片内存占用钉死。
  • maxFileSize=10GB 是单文件硬上限;配额层面的"用户总空间够不够"由 CheckStorageAvailable 另行负责(见一致性小节)。

⚠️ :钳制用的是"夹取"而不是"拒绝"。如果客户端传 chunkSize=1,直接拒绝会让前端体验差;夹到 minChunkSize 既能保安全又不打断上传。但分片数仍要在 InitUpload 之后校验 ChunksTotal 不超过某个合理上限,否则 10GB 文件被夹到 1MB 上限会产生 1 万次 UploadPart,仍要给并发和 DB 写入兜底(见第九节信号量)。

再补一个工程建议:钳制 chunkSize 之后,服务端应顺手校验 ChunksTotal 不超过一个硬上限(比如 1 万或 2 万)。因为即便 chunkSize 被夹到 1MB,10GB 文件仍会产生约 1 万个分片,每个分片一次 UploadPart + 一条 upload_chunks 记录。ChunksTotal 过大会带来两个副作用:一是 DB 里 upload_chunks 表瞬间多 1 万行,合并时 ListChunksBySession 一次查 1 万行有压力;二是客户端要发 1 万次请求,单纯网络往返就不可忽略。所以更优雅的做法是:ChunksTotal 超过阈值时,直接把 chunkSize 进一步放大(在 [minChunkSize, maxChunkSize] 内取满足上限的最小值),把分片数压到合理区间。本项目用 1MB 下限是出于"兼顾细粒度续传"的考虑,生产里也可以按"文件大小 → 推荐分片数 100~1000"动态选 chunkSize,这是更贴近大厂实现的做法。

Q7. 每片上传时,服务端还要做哪些校验?

答:

类比:收快递时,你不仅看箱子编号在 1~12 之间,还要看箱子是不是空的、是不是明显超重。编号越界、空箱、超重的箱子都不该收。

工程内容FileUsecase.UploadPart 在真正写盘前做三道校验:

// internal/biz/file.go
// 1) 分片序号必须在 [1, ChunksTotal],防越界 / 写错位置
if partNumber < 1 || partNumber > int(session.ChunksTotal) {
    return nil, perrors.BadRequest("INVALID_PART_NUMBER", "分片序号越界")
}
// 2) 分片不能为空,防空分片污染合并结果
if len(data) == 0 {
    return nil, perrors.BadRequest("EMPTY_PART", "分片数据为空")
}
// 3) 单分片不能超过上限,防超大分片撑爆存储/内存
if int64(len(data)) > maxChunkSize {
    return nil, perrors.BadRequest("PART_TOO_LARGE", "分片大小超过系统上限")
}

再加上 storage 层的 os.Stat(partPath) 幂等保护:同一片重复传会返回 biz.ErrPartAlreadyUploaded,让客户端可以安全重试而不产生重复文件。biz 层还要先校验 session.Status == 0(进行中)和 session.UserID == userID(越权保护),否则已结束或别人的会话也能被写入。


五、断点续传:网络断了接着传

Q8. 断点续传是怎么实现的?客户端怎么知道该补传哪几片?

答:

类比:你下载 12 箱冰箱零件,收到第 5 箱时快递员摔了腿(断网)。等你康复,不需要把 12 箱全重寄——你打电话问"我目前收到了哪几箱?“对方说"1、2、3、4”,你只需补寄 5~12。断点续传就是"记住已到货的箱子编号,只补缺失的"。

工程内容:本项目用"相同 hash + 相同 size 的未完成会话"来识别"这就是同一个文件的上一次上传"。

1) 恢复会话InitUpload 里,当客户端带了 hash 时,先用 FindByHash 看是否秒传;若不秒传,再用 FindByHashAndSize 找"同 hash 同大小、且状态为进行中(0)或已暂停(3)“的历史会话:

// internal/biz/file.go(节选)
if hash != "" {
    existing, _ := uc.fileRepo.FindByHash(ctx, userID, hash)
    if existing != nil && existing.Status == 0 {
        return &UploadSession{Status: 1}, nil // 已存在同 hash 文件 → 直接秒传
    }
    // 断点续传:查找同 hash + 同大小的进行中或已暂停会话
    if uc.uploadRepo != nil {
        resumed, _ := uc.uploadRepo.FindByHashAndSize(ctx, userID, hash, fileSize)
        if resumed != nil && (resumed.Status == 0 || resumed.Status == 3) {
            if resumed.Status == 3 {
                _ = uc.uploadRepo.ResumeSession(ctx, resumed.UploadID) // 已暂停 → 恢复进行中
                resumed.Status = 0
            }
            return resumed, nil // 直接复用旧会话
        }
    }
}

对应的仓储查询(internal/data/upload.go)只挑 status IN (0, 3) 的会话,按最新创建时间取一条:

// internal/data/upload.go
func (r *uploadRepo) FindByHashAndSize(ctx context.Context, userID uint64, hash string, fileSize int64) (*biz.UploadSession, error) {
    var po model.UploadSession
    if err := r.db.WithContext(ctx).
        Where("user_id = ? AND file_hash = ? AND file_size = ? AND status IN (0, 3)", userID, hash, fileSize).
        Order("created_at DESC").
        First(&po).Error; err != nil {
        // ...
    }
    return toBizUploadSession(&po), nil
}

2) 查已传分片:客户端调用 CheckParts,服务端返回"已经收到哪些分片号”,客户端据此计算差集补传:

// internal/biz/file.go
func (uc *FileUsecase) CheckParts(ctx context.Context, userID uint64, uploadID string) ([]int32, error) {
    session, err := uc.uploadRepo.FindSessionByUploadID(ctx, uploadID)
    if err != nil {
        return nil, err
    }
    if session.UserID != userID {
        return nil, ErrForbidden
    }
    return uc.uploadRepo.ListReceivedPartNumbers(ctx, session.ID)
}

ListReceivedPartNumbersupload_chunks 表把所有 chunk_index 查出来(internal/data/upload.go)。注意:DB 里的分片记录 + 磁盘里的 part_%05d 应该保持一致;本项目以 DB 的 chunks_received 计数和 upload_chunks 表为准来判定是否齐全。

⚠️ :断点续传判定只靠 hash + size 存在理论碰撞风险(两个不同文件恰好同 hash 同大小概率极低但非零),且要求客户端"同一文件用同一 hash"。工程上更稳妥是再带一个客户端生成的 fileID 或客户端本地持久化 uploadID,本项目用 hash+size 作为启发式恢复,已能覆盖绝大多数"刷新页面/断网重连"场景。

断点续传流程如下:

flowchart TD
    A["客户端 断网后重连
带 hash + size 重新 InitUpload"] --> B["biz 先 FindByHash 秒传?"] B -->|命中| C["直接秒传"] B -->|未命中| D["FindByHashAndSize
找同 hash+size 的 0/3 会话"] D -->|找到| E["ResumeSession 恢复
复用旧 uploadID"] D -->|没找到| F["新建会话"] E --> G["CheckParts 查已收分片号"] F --> G G --> H["客户端算差集
只补传缺失分片"] H --> I["补传 UploadPart"] I --> J["MergeParts 合并"]

客户端如何用 CheckParts 的结果补传?核心是一行"差集"计算。假设服务端返回已收分片号 [1,2,4,5],总片数 ChunksTotal=6,那么需要补传的就是 [3,6]

// 客户端伪代码
received := checkParts(uploadID)            // 如 [1,2,4,5]
need := []int{}
for i := 1; i <= chunksTotal; i++ {
    if !slices.Contains(received, int32(i)) {
        need = append(need, i)              // [3,6]
    }
}
for _, n := range need {
    uploadPart(uploadID, n, readChunk(n))   // 只补缺失片
}

这里有个容易忽略的点:ChunksTotal 必须以服务端钳制后的 chunkSize 为准,而不是客户端最初的 chunkSize。因为服务端可能把客户端传来的 chunkSize 夹到了 [minChunkSize, maxChunkSize],导致实际分片数和客户端预期不一致。如果客户端用"未钳制"的分片数去算差集,会出现"少传最后一片"或"多传不存在的片"的错位。本项目 InitUpload 返回 session.ChunksTotal,客户端应直接用这个值。


六、流式下载、预览与物理/逻辑一致性

Q9. 下载为什么要返回 io.ReadCloser?预览的 URL 怎么来?

答:

类比:看一部 4GB 电影,你不会先把它完整下载到手机内存再播放,而是"边下边播"(流式)。如果非要把整部电影读进内存,手机直接爆了。io.ReadCloser 就是"边读边传的水管",数据从磁盘流向 HTTP 响应体,全程不落地到内存。

工程内容

1) 流式下载 DownloadlocalStorage.Download 直接 os.Open 返回 *os.File,它本身就是 io.ReadCloser

// internal/data/storage/local.go
func (s *localStorage) Download(ctx context.Context, storagePath string) (io.ReadCloser, error) {
    f, err := os.Open(storagePath)
    if err != nil {
        if os.IsNotExist(err) {
            return nil, biz.ErrStorageFileNotFound
        }
        return nil, fmt.Errorf("storage: failed to open file: %w", err)
    }
    return f, nil
}

biz.GetFileStream 把它透传给 service 层,由 HTTP 框架把 reader 管道接到响应体:

// internal/biz/file.go
func (uc *FileUsecase) GetFileStream(ctx context.Context, userID uint64, fileID uint64) (*File, io.ReadCloser, string, error) {
    file, err := uc.fileRepo.FindByID(ctx, fileID)
    // ... 越权 + 回收站校验 ...
    reader, err := uc.storage.Download(ctx, file.Path)
    return file, reader, file.Type, nil
}

注意 biz.Download(小文件版)里用了 io.ReadAll 把内容整块读进 []byte 返回——那是给"需要完整字节做处理"的场景(如生成缩略图)。大文件预览/下载请务必走 GetFileStream 的流式路径,否则 4GB 文件直接 OOM。

2) 预览 URL GetFileURL

  • 本地存储返回 file:// + 绝对路径(localStorage.GetFileURL)。
  • MinIO 用 PresignedGetObject 生成限时 1 小时的签名 URL(Preview 里传 1*time.Hour),客户端拿这个 URL 直接在浏览器看,不用把流量过我们的服务:
// internal/biz/file.go
func (uc *FileUsecase) Preview(ctx context.Context, userID uint64, fileID uint64) (string, string, error) {
    // ... 越权 + 回收站校验 ...
    previewType := guessPreviewType(file.Type)
    url, err := uc.storage.GetFileURL(ctx, file.Path, 1*time.Hour) // 限时 1h
    return url, previewType, err
}
// internal/data/storage/minio.go
func (s *minioStorage) GetFileURL(ctx context.Context, storagePath string, expire time.Duration) (string, error) {
    presignedURL, err := s.client.PresignedGetObject(ctx, s.bucket, storagePath, expire, nil)
    if err != nil {
        return "", fmt.Errorf("storage: minio 生成预签名 URL 失败: %w", err)
    }
    return presignedURL.String(), nil
}

MinIO 的分片清理走后台 go s.cleanupChunks(...)(合并成功后异步删分片对象),避免阻塞合并响应——这是与本地存储"同步 RemoveAll“的一个实现差异。

再强调一次 Download(整块读 []byte)与 GetFileStream(流式 io.ReadCloser)的取舍:biz.Downloadio.ReadAll 把内容完整读进内存再返回,适合"上传后立刻需要字节做处理”(如服务端校验、生成缩略图)的小文件;而 GetFileStream 透传 reader,由 HTTP 框架把水管接到响应体,内存恒定。面试被问"大文件下载怎么做"时,标准答案就是后者——永远不要让完整文件体经过应用进程的内存。更进一步,配合 HTTP Range 请求还能实现"断点下载"(见面试延伸第 3 点),让前端拖进度条秒开、断网从断点续下。

流式下载与预览的两条路径:

flowchart TD
    A["客户端请求文件"] --> B{"大文件 还是 预览?"}
    B -->|大文件下载| C["GetFileStream
storage.Download 返回 io.ReadCloser"] C --> D["HTTP 边读边传
零拷贝 不占内存"] B -->|浏览器预览| E["Preview
storage.GetFileURL"] E --> F["本地: file:// 绝对路径"] E --> G["MinIO: PresignedGetObject
限时 1h 签名 URL"] F --> H["前端直接打开"] G --> H

Q10. 物理文件与 DB 记录怎么保持一致?删除顺序有讲究吗?

答:

类比:图书馆里"书本实体"和"借阅卡上的编号"必须对应。丢书(物理删了但卡还在)→ 读者按卡去找,扑空;卡丢了书还在(卡删了物理还在)→ 书成了没人认领的孤儿。我们要么"先销卡再下架书",要么"下架书失败就保留卡"。

工程内容:本项目在两处严守一致性:

  • 上传/合并成功、DB 落库失败 → 删物理文件回滚(前面小文件与 MergeParts 都体现了)。MergeParts 里有三处回滚点:UUID 生成失败删已合并文件、DB Create 失败删已合并文件。
  • 删除文件:先确保物理删除成功,再标记 DB(或反过来以 DB 为"真相源")。本项目的 SecUpload / Copy 等"逻辑新增"都配了 AddUsedStorage 记账;删除虽不在本次必读文件里,但设计原则是物理删除成功后才改 DB 状态,因为物理文件一旦误删不可恢复,而 DB 状态可以再修。

另一个一致性细节是存储配额用分布式锁保护internal/biz/storage.goLocker 接口 + updateStorageWithLock):并发上传时多个请求同时 AddUsedStorage,如果不加锁,“读取旧值 → +delta → 写回"会成为竞态,导致配额计数偏小。updateStorageWithLock 在更新前后加分布式锁,并在更新后清缓存,保证配额准确。这属于"分布式锁篇"的内容,本文点到为止。

⚠️ 千万不要先删 DB 再删物理。如果 DB 先删、物理删除时磁盘满了/权限错了失败,这份文件就永远成了"DB 不认、磁盘还在"的孤儿,且无法再通过正常流程发现与回收。正确顺序是"物理删除成功 → 再改 DB”,或者"DB 标记删除 → 异步任务确认物理删除"。


七、MIME 探测与安全拦截

Q11. 怎么知道一个文件是什么类型?危险文件怎么拦?

答:

类比:你收快递,看箱子外观(扩展名)猜里面是书还是炸弹。但外观会骗人——于是你再看"快递单上的品类标注"(系统 MIME 数据库)。如果连快递公司都没登记这个品类(精简系统没有 mime.types),你就掏出自己手写的"常见品类速查表"(内置映射)兜底。对于明显是炸药的品类(.exe / .sh),直接拒收。

工程内容detectMimeType 三级兜底:

// internal/biz/file.go
var builtinMimeTypes = map[string]string{
    ".mp4":  "video/mp4",
    ".mov":  "video/quicktime",
    // ... 视频/音频/图片/文档/压缩包 共数十项 ...
    ".zip":  "application/zip",
    ".gz":   "application/gzip",
}

func detectMimeType(filename string) string {
    ext := filepath.Ext(filename)
    if ext == "" {
        return "application/octet-stream"
    }
    // 1) 优先系统 MIME 数据库
    if mimeType := mime.TypeByExtension(ext); mimeType != "" {
        return mimeType
    }
    // 2) 回退内置映射(Docker 精简镜像常缺 mime.types)
    if mimeType, ok := builtinMimeTypes[strings.ToLower(ext)]; ok {
        return mimeType
    }
    return "application/octet-stream"
}

为什么要内置 builtinMimeTypes:Go 的 mime.TypeByExtension 底层读系统的 /etc/mime.typesDocker 用 scratch / distroless 精简镜像时这个文件常常不存在,结果 .mp4 视频被识别成 application/octet-stream,前端 <video> 标签直接播不了。内置映射补上最常见的几十种,保证核心格式在任意环境都识别正确。

再深入一层:mime.TypeByExtension 并不是每次调用都去读磁盘,而是进程启动时把系统的 MIME 数据库(一份"扩展名 → Content-Type"的映射表)一次性读进内存缓存,之后查的是内存。因此它的可靠性完全取决于基础镜像里有没有那份文件——scratch / distroless / 部分 alpine 精简镜像往往没有,于是所有扩展名都查不到,video/mp4 退化成 application/octet-stream。本项目把最常见几十种格式硬编码进 builtinMimeTypes,本质是给"系统库缺失"留一道应用层兜底:先问系统库,系统库空再查自己的小表,两路都没有才退回 octet-stream。这种"三级兜底"的写法特别能体现生产环境意识,面试讲出来很有分量——它会让面试官觉得你真的在生产容器里踩过这个坑。

危险类型拦截用黑名单 dangerousExtensions + isDangerousFileType,在 Upload / SecUpload / InitUpload 之前都先过一遍:

// internal/biz/file.go
var dangerousExtensions = map[string]bool{
    ".exe": true, ".sh": true, ".php": true, ".jsp": true,
    ".py": true, ".dll": true, ".app": true, ".jar": true,
    // ... 共 28 种可执行/脚本类扩展名 ...
}

func isDangerousFileType(fileName string) bool {
    ext := strings.ToLower(filepath.Ext(fileName))
    return dangerousExtensions[ext]
}

此外 sanitizeFileName 会把 ../..\\<>:|?* 等危险字符剥离,并限制长度 ≤255,防止路径穿越与文件系统非法名。

⚠️ 扩展名黑名单只是第一道闸,不是安全终点.php 被拦了,攻击者可以把 webshell 改名 .phtml 或利用服务器配置缺陷。真正安全还要结合"上传目录禁止执行 CGI/脚本"“存储与 Web 根分离"“对图片做二次渲染"等。面试里能说出"黑名单不够、还需执行隔离"会很加分。


八、孤儿分片清理

Q12. 服务崩溃或客户端断连,残留的分片目录怎么处理?

答:

类比:快递分拣中心偶尔有"寄件人付了首重就失联"的半截包裹,占着货架。保洁员每天巡场,把"超过 N 天没动静"的货架整排清掉。孤儿分片就是这些"没人来续传也没人来取消"的 chunks/<uploadID> 目录。

工程内容Storage 接口的 CleanupOrphanChunks(ctx, maxAge) 负责扫描 chunks/ 目录,删掉修改时间早于 maxAge 的子目录。本地实现:

// internal/data/storage/local.go
func (s *localStorage) CleanupOrphanChunks(ctx context.Context, maxAge time.Duration) ([]string, error) {
    chunksRoot := filepath.Join(s.uploadDir, "chunks")
    entries, err := os.ReadDir(chunksRoot)
    if err != nil {
        if os.IsNotExist(err) {
            return nil, nil // 没有 chunks 目录 = 没有孤儿
        }
        return nil, fmt.Errorf(...)
    }
    now := time.Now()
    removed := make([]string, 0, len(entries))
    for _, entry := range entries {
        if ctx.Err() != nil {
            return removed, ctx.Err()
        }
        if !entry.IsDir() {
            continue
        }
        uploadID := entry.Name()
        info, err := entry.Info()
        if err != nil {
            continue // 取不到信息不删,宁可漏清也不误删
        }
        // 用目录 ModTime 判断;超过 maxAge 视为孤儿
        if now.Sub(info.ModTime()) > maxAge {
            if err := os.RemoveAll(filepath.Join(chunksRoot, uploadID)); err != nil {
                fmt.Printf("storage: failed to remove orphan chunk dir %s: %v\n", uploadID, err)
                continue
            }
            removed = append(removed, uploadID)
        }
    }
    return removed, nil
}

MinIO 版则是 ListObjectschunks/ 前缀,按 LastModified 过滤后批量 RemoveObjects

两个安全细节值得记:

  1. 用目录的 ModTime 判断:每次 UploadPart 写文件都会更新目录 mtime,所以"还在续传"的会话 mtime 是新的,不会被误清。只有彻底静止超过 maxAge 的才是真孤儿。
  2. entry.Info() 取不到就跳过:清孤儿是"尽力而为”,宁可漏清(下次定时任务再清)也不能因为一次 Stat 失败就 RemoveAll 误删一个活跃会话。
  3. ctx.Err() 检查:超时被取消时及时返回已清理列表,避免长时间占用。

清理定时任务由外部调度(cron / 启动 goroutine)周期性调用,返回的 []string 可用于打点监控"本次清理了多少孤儿”。

flowchart TD
    A["定时任务 调用 CleanupOrphanChunks"] --> B["os.ReadDir chunks/"]
    B --> C{"每个子目录"}
    C --> D["取目录 ModTime"]
    D --> E{"静止时长 超过 maxAge?"}
    E -->|否 还在续传| F["保留"]
    E -->|是 孤儿| G["os.RemoveAll 删除"]
    G --> H["加入 removed 列表"]
    F --> I["返回 removed 用于监控"]
    H --> I

九、上传并发限流

Q13. 高并发上传时,怎么保护磁盘和数据库?

答:

类比:收费站只有 50 个车道,车多了就排队,排太久(10 秒)就劝返"请稍后再来",而不是让 5000 辆车同时挤上高速把路压垮。uploadSemaphore 就是这个"车道数有限的收费站"。

工程内容FileUsecase 持有一个带缓冲的 channel 作为信号量,uploadSemaphoreSize = 50(压测得出并发 50 时吞吐峰值约 190 QPS,再高反而倒退):

// internal/biz/file.go
const uploadSemaphoreSize = 50

type FileUsecase struct {
    // ...
    uploadSemaphore chan struct{} // 上传并发限流信号量
}

func NewFileUsecase(...) *FileUsecase {
    return &FileUsecase{
        // ...
        uploadSemaphore: make(chan struct{}, uploadSemaphoreSize),
    }
}

UploadPart 里用 select 抢信号量,超时 10 秒返回 429 UPLOAD_BUSY

// internal/biz/file.go
semTimeout := 10 * time.Second
select {
case uc.uploadSemaphore <- struct{}{}:
    defer func() { <-uc.uploadSemaphore }()
case <-time.After(semTimeout):
    return nil, ErrUploadBusy // 429,前端友好提示"稍后重试"
case <-ctx.Done():
    return nil, ctx.Err()
}

为什么用"信号量 + 超时"而不是"无限制接":分片上传每一步都要写磁盘(io.Copy)+ 写 DB(CreateChunk / IncrementChunks)。并发过高时,磁盘 I/O 争抢和 MySQL 行锁会让整体吞吐不升反降,且 P95 延迟暴涨。限流把并发钉在拐点附近,既保吞吐又保延迟。10 秒是权衡值:短了容易误伤正常排队,长了拉高 P95。

⚠️ :信号量是进程内的,多副本部署时每个 Pod 各自有 50 个车道,全局实际并发 = 50 × Pod 数。如果要全局精确限流,需要 Redis 令牌桶之类的分布式限流。面试能点出"进程内信号量 vs 分布式限流"的差异,是加分项。


十、面试延伸

这部分把本文知识点往"大厂真实架构"推一步,面试常追问。

1) 大文件直传 OSS 的 presigned URL 方案

本项目 MinIO 的 GetFileURL 已经用了 PresignedGetObject。上传也可以对称地用 presigned PUT/POST URL:客户端先向我们的服务"申请一个上传凭证"(带过期时间和对象 key 的签名 URL),然后客户端绕过我们的服务,直接把文件 PUT 到 OSS。这样我们服务不扛文件流量,只管"发凭证 + 收完成回调 + 写 DB",能轻松支撑 TB 级文件。分片则结合 OSS/MinIO 的服务端 InitiateMultipartUpload / UploadPart / Complete,由客户端逐片直传,服务端只在最后 Complete 时校验 ETAG 列表。

2) 分片合并的原子性

本地实现里,合并是"边读边写最终文件 → Sync → 删分片目录"。问题在于:如果 Sync 成功但 os.RemoveAll 之前进程崩了,会留下"最终文件 + 分片目录"两份。由于 finalPath = chunks 之外的日期目录 + uploadID.merged,且 DB 的 File.Path 指向 finalPath只要 DB 没标记完成,这次合并就不算成功,客户端重传会重新 InitUpload 生成新的 uploadID,旧的分片目录和 .merged 都由 CleanupOrphanChunks 兜底清掉。换句话说,“DB 会话状态"才是合并是否成功的真相源,物理中间产物靠孤儿清理兜底,这是最终一致性的典型思路。

3) 如何支持断点下载(HTTP Range 请求)

断点续传是"上传"侧;“下载"侧对应 HTTP Range。实现要点:Download 返回的 io.ReadCloser 换成支持 Seek*os.File,service 层解析 Range: bytes=0-1023file.Seek(offset, io.SeekStart) 后只把对应区间流式写回,并带 206 Partial Content + Content-Range + Accept-Ranges: bytes 响应头。这样几十 GB 的视频也能"拖进度条秒开"“断网从断点接着下”。对象存储(OSS/S3)原生支持 Range GET,直接传 Range 头即可,无需自己 Seek

  1. 秒传 / 分片与对象存储去重的边界:当你把存储换成 OSS,秒传其实可以下沉到对象存储层——OSS 的 PutObject 支持按 Content-MD5 或客户端的 x-oss-copy-source 做服务端去重,但更常见的仍是"应用层按 SHA-256 去重 + 复用 Path"。要提醒的是:对象存储按"对象键"计费的模型下,秒传省的是"上传流量 + 存储字节”,但"多个逻辑文件指向同一对象键"依然要求应用层做引用计数,否则一个 DeleteObject 会让所有秒传记录失效。这与本地磁盘的引用计数思路完全一致,只是操作对象从"文件"变成"对象键”。

  2. 可观测性:分片上传这类长链路操作一定要打点——InitUpload 成功率、UploadPart 平均大小与耗时、MergeParts 耗时分布、CleanupOrphanChunks 每次清理的孤儿数。这些指标能在"上传变慢"“磁盘莫名涨满"时第一时间定位是客户端分片太小、还是孤儿清理没跑、还是合并 Sync 成了瓶颈。


自测题与动手练习

自测题(5 道)

  1. 本项目 Storage 接口有哪 9 个方法?分别服务于什么场景?为什么 biz 层要依赖接口而非具体实现?
  2. var _ biz.Storage = (*localStorage)(nil) 这行代码有什么用?删掉它会怎样?OSS 为什么只写 stub 而不是报错 panic?
  3. 秒传为什么用 SHA-256 判断"同一文件”?只用 hash 够不够?本项目 FindByHashAndSize 相比 FindByHash 多了什么约束,解决了什么问题?
  4. MergeParts 合并分片时,为什么 sort 排序、bufio.NewWriterSize(dest, 1024*1024) 1MB 缓冲、dest.Sync() 三件事缺一不可?各自解决什么?
  5. 分片大小钳制 minChunkSize=1MB / maxChunkSize=100MB / maxFileSize=10GB 分别防什么攻击或故障?如果不钳制,客户端传 chunkSize=1 会发生什么?

动手练习(3 件)

  1. 在本地起 MinIO(docker 一行命令),把项目的 NewLocalStorage 换成 NewMinioStorage,上传一个文件后用 mc 命令行验证对象确实落在 年/月/日/ 前缀下,并对比本地存储与 MinIO 在 CompleteMultipartUpload 后清理分片的方式差异(本地同步 RemoveAll vs MinIO 后台 go cleanupChunks)。
  2. 写一个单测:并发 200 个 goroutine 同时调 UploadPart,用 go test -race 验证 uploadSemaphore 确实把并发钳制在 50,且不会出现 ChunksReceived 计数少于实际分片数的情况(提示:看 IncrementChunksgorm.Expr("chunks_received + 1") 原子更新)。
  3. ossStorageUploadGetFileURL 补上真实 aliyun-oss-go-sdk 实现(参考 minio.go 的 presigned 思路用 bucket.SignURL),跑通后在 CleanupOrphanChunks 里用 ListObjectsV2 + DeleteObjects 实现孤儿清理,验证编译期 var _ biz.Storage 仍通过。

本章小结

如果在面试里被要求"从头讲一遍文件上传",建议用下面这条主线组织语言,比背八股更有说服力:先讲抽象(为什么要有 Storage 接口、编译期 var _ biz.Storage 校验、OSS 用 stub 不阻塞交付)→ 再讲小文件io.Copy 流式 + 按日期分目录,避免 OOM 和单目录爆炸)→ 接着讲优化(秒传 FindByHash 省存储、FindByHashAndSize 恢复断点续传会话)→ 然后讲大文件的分片三步走InitUploadchunks/<uploadID>UploadPartpart_%05dMergeParts 排序 + 1MB bufio + Sync 落盘)→ 再补安全闸门(分片大小钳制 [1MB,100MB]、单文件 ≤10GB、detectMimeType 三级兜底、isDangerousFileType 黑名单)→ 最后讲稳定与清理(上传成功写 DB 失败回滚删物理、删除先物理后逻辑、CleanupOrphanChunks 清孤儿、uploadSemaphore 限流 50)。这条线把"设计、性能、安全、一致性、可运维"五个维度全串起来了,面试官要的往往就是这种完整工程观,而不是某个 API 怎么调。

本文以一个真实的 Kratos 云盘项目为蓝本,把"对象存储与文件分片上传/下载"这条面试高频线串了起来:

  • 抽象是地基biz.Storage 接口把业务与本地磁盘 / MinIO / OSS 解耦,var _ biz.Storage = (*X)(nil) 做编译期实现校验,OSS 用 stub 留待接入而不阻塞编译。
  • 小文件走 io.Copy 流式 + 按日期分目录,避免单目录文件爆炸;上传成功写 DB 失败要回滚删物理文件。
  • 秒传FindByHash 复用 Path,但配额仍要按 Size 记账;断点续传FindByHashAndSize 恢复会话、CheckParts 告知已收分片。
  • 分片上传三步走:InitUploadchunks/<uploadID>(biz 与 storage 同名防孤儿)→ UploadPartpart_%05d(幂等、校验序号/非空/上限)→ MergeParts 排序 + 1MB bufio 合并 + Sync 落盘 + 删分片目录。
  • 安全闸门:分片大小钳制 [1MB,100MB]、单文件 ≤10GB;detectMimeType 三级兜底(系统库 → 内置 builtinMimeTypesoctet-stream);isDangerousFileType 黑名单拦截可执行脚本。
  • 资源清理与稳定CleanupOrphanChunks 按 mtime 清静止超期的孤儿分片;uploadSemaphore(50)限流保护磁盘与 DB,超时返回 429
  • 下载零拷贝Download / GetFileStream 返回 io.ReadCloser 边读边传;Previewfile://(本地)或限时 1h 的 presigned URL(MinIO/OSS)。

把这些点连成一条链路,你在面试里讲"文件上传"就不再是背八股,而是能从一个真实项目讲出抽象设计、边界校验、一致性、限流、清理的完整工程权衡——这正是资深后端区别于初级的关键。

复习提示:
  • Storage 接口抽象的价值:不只是一行代码,而是让 OSS/MinIO/本地磁盘可以热替换,扩容时只需改 data 层配置,biz 层完全无感。
  • 分片上传三步走的精髓:Init 建目录 → UploadPart 幂等写分片(用序号命名防重复)→ MergeParts 排序合并——每个步骤都有明确的失败点和补偿机制。
  • 安全门三层防护:文件名清洗(防路径穿越)→ 类型检测三级兜底(系统库→内置白名单→octet-stream)→ 执行文件黑名单,任何一层漏掉都有另一层补位。
  • 孤儿分片清理:按 mtime 而非 createdAt 判断,因为用户上传可能跨天;cleanup 必须定期执行,否则分片目录会无限膨胀。
  • 下一篇讲缓存设计——它是文件查询的性能关键,配合索引和分页形成完整的数据访问层方案。
面试官
断点续传里,客户端如何知道自己已经上传了哪些分片?服务器端怎么做到幂等?
候选人

客户端识别已上传分片
调用 CheckParts 接口,服务端返回该 uploadID 下已存在的分片列表(通过扫描 chunks/<uploadID>/ 目录或查 DB)。客户端对比本地分片和已上传分片,只重传缺失的。

服务端幂等保证
① 分片文件名用序号命名 part_%05d(不足 5 位补零),覆盖写入时内容相同则结果一致
② 每次 UploadPart 都校验:

  • 序号是否合法(1~MaxParts)
  • 文件大小是否 >0
  • 分片大小是否在 [1MB, 100MB] 范围内

    面试加分点:提到"覆盖写入虽然简单但不够安全"——如果同一分片被恶意重复提交大文件,会浪费磁盘。所以项目中还做了配额预扣减(先 ReserveStorage 再上传)。
About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!