学习目标
学完本章你应该能够:
- 说清楚
biz.Storage接口如何把"业务要做什么"和"底层怎么做"解耦,并用编译期断言var _ biz.Storage = (*localStorage)(nil)保证实现不漏方法。 - 讲清本地磁盘存储与 MinIO 对象存储在路径生成、分片合并上的本质差异(本地
io.CopyBuffervs MinIOComposeObject服务端合并)。 - 解释为什么按
年/月/日分目录,以及为什么单目录文件过多会拖垮文件系统。 - 描述分片上传的完整生命周期,以及"孤儿分片"是怎么产生的、怎么清理。
- 讲清配额(quota)这条线:
used_storage如何被原子增减、CheckStorageAvailable如何在上传前预检、缓存GetStorageInfo的穿透/击穿/雪崩与失效时机、以及并发上传下"为什么光靠分布式锁还不够"。 - 在面试中把"抽象先行、实现可替换"与"配额原子更新 + DB 上界是唯一裁判"讲成两段有结构的话。
前置知识:
- Go 语言基础(
interface、编译期断言、io.Reader/Writer、context) - 文件系统与对象存储的基本概念(本地路径 vs S3 对象 key)
- MySQL 行锁、
gorm.Expr原子表达式、Redis 基础命令(GET/SET/DEL/TTL) - Kratos 项目结构(
biz业务层与data数据层分层的依赖倒置)
本章你会动手做的事:
- 画出
Storage接口 9 个方法的职责草图,并标注哪些是分片上传生命周期。 - 对比本地与 MinIO 的
CompleteMultipartUpload,说出"服务端合并"省掉了哪部分带宽。 - 给
UpdateUsedStorageAtomic的增加分支补上WHERE used_storage + ? <= total_storage条件,并画一张"并发上传是否超卖"的判定流程图。
一、技术栈与中间件
| 技术 / 中间件 | 用途说明 |
|---|---|
Go 标准库 os | 文件 / 目录的创建、打开、删除(os.Create、os.Open、os.MkdirAll、os.Remove、os.ReadDir、os.RemoveAll) |
Go 标准库 io | 流式读写与零拷贝复制(io.Copy、io.CopyBuffer、io.Reader / io.ReadCloser) |
Go 标准库 bufio | 带缓冲的 Writer,用于大文件分片合并时减少 write 系统调用次数 |
Go 标准库 path/filepath | 跨平台路径拼接(filepath.Join、filepath.Abs) |
Go 标准库 sort | 分片合并前按 PartNumber 排序,确保最终文件字节顺序正确 |
Go 标准库 time | 按年/月/日生成存储子目录、孤儿分片过期判断 |
Go 标准库 log/slog | 结构化日志,用于 MinIO 初始化、孤儿分片清理告警 |
github.com/minio/minio-go/v7 | MinIO / S3 兼容对象存储 SDK,提供 PutObject、GetObject、ComposeObject、PresignedGetObject、RemoveObjects 等接口 |
github.com/google/wire | 依赖注入框架,通过 ProviderSet 注册存储工厂 NewStorage |
github.com/aliyun/aliyun-oss-go-sdk/v3(计划中) | 阿里云 OSS SDK,目前以 stub 形式预留扩展点 |
| GORM | 数据访问层(data),gorm.Expr("used_storage + ?") 做数据库内原子增减;Where("id=? AND used_storage >= ?") 做扣减防超卖 |
| MySQL | 用户表(model.User)持久化,依赖行锁 + 条件更新保证 used_storage 并发一致;total_storage 字段为有效配额 |
| Redis(通过 biz.Cache 接口) | 存储统计信息缓存(GetStorageInfo)、分布式锁(Locker)、缓存失效 |
Kratos perrors | 业务错误码体系,例如 ErrStorageInsufficient(413) |
| Storage 抽象接口设计 | biz.Storage 接口屏蔽底层差异,业务层只依赖抽象 |
| 按年/月/日分目录存储 | datePath() / dateKey() 生成 2026/08/04/ 路径,避免单目录文件过多 |
分布式锁 biz.Locker | 用户存储配额更新时的可选并发串行化(注意:非正确性唯一保证) |
二、实现思路流程(总体)
类比:存储模块就像一家快递公司的"分拣中心"——上层业务只管把包裹交给前台(
Storage接口),至于包裹走本地仓还是走云仓(MinIO/OSS),业务完全不关心。换仓库只换后台,前台流程不变。而"已用空间"像公司的账本,每一笔入库/出库都在账本上当场加减,谁也插不了队。
flowchart TB
B[biz 业务层] -->|只依赖接口| I[Storage 接口
9 个方法]
I -->|local 驱动| L[localStorage
本地磁盘]
I -->|minio 驱动| M[minioStorage
S3 对象存储]
I -->|oss 预留| O[ossStorage stub
扩展点]
L --> F[Wire 工厂 NewStorage
按 Driver 选择实现]
M --> F
U[FileUsecase / RecycleUsecase] -->|上传前预检| Q[CheckStorageAvailable
读缓存判断余量]
U -->|写盘后记账| A[AddUsedStorage / SubUsedStorage]
A -->|分布式锁| Lk[updateStorageWithLock]
Lk -->|原子 SQL| DB[(MySQL users 表
used_storage 增减)]
Lk -->|失效| C[(Redis 存储统计缓存)]这张图说明两件事:业务层永远只看到 Storage 接口,具体实现由工厂在启动时决定;同时"上传/删除"会顺带走一条配额记账链路——预检、原子扣减、缓存失效。存储模块整体遵循「抽象先行,实现可替换」的设计哲学,配额这条线则遵循「预检兜底 UX,DB 上界兜底正确性」。
存储模块整体实现思路按以下流程展开:
抽象接口设计:在
internal/biz/storage.go中定义Storage接口,包含Upload、Download、Delete、InitMultipartUpload、UploadPart、CompleteMultipartUpload、AbortMultipartUpload、CleanupOrphanChunks、GetFileURL共 9 个方法。业务层(biz)只依赖此接口,完全不感知底层是本地磁盘、MinIO 还是 OSS。本地存储实现:在
internal/data/storage/local.go中实现localStorage,使用 Go 标准库os/io操作文件系统。上传时通过datePath()在uploadDir下按年/月/日自动创建子目录,避免单目录文件爆炸。MinIO 实现:在
internal/data/storage/minio.go中实现minioStorage,通过minio-goSDK 与 S3 协议兼容的对象存储交互。分片合并使用ComposeObjectAPI 在服务端完成,无需下载到本地再拼接。OSS 扩展规范:在
internal/data/storage/oss_stub.go中以 stub(桩)形式预留阿里云 OSS 实现,所有方法返回ErrStorageNotImplemented,并通过注释明确每个方法对应的 OSS SDK 调用,便于后续按规范补全。工厂选择:在
internal/data/storage/storage.go中通过NewStorage根据conf.Storage.Driver字段(local/minio)选择具体实现,并集成到 Wire 的ProviderSet中。存储路径按日期分目录:本地存储用
filepath.Join(uploadDir, "2026", "08", "04"),MinIO 用2026/08/04/对象前缀,使文件天然分散。孤儿分片清理:分片上传过程中若服务崩溃或客户端断开,
chunks/{uploadID}/下会残留分片。CleanupOrphanChunks扫描该目录,删除ModTime早于maxAge的子目录,回收磁盘 / 对象存储空间。配额记账:上传/分享转存/复制/分片完成前先
CheckStorageAvailable预检余量,落盘并写库成功后AddUsedStorage原子增加used_storage;回收站永久删除时SubUsedStorage原子减少;每次增减后updateStorageWithLock失效 Redis 缓存。GetStorageInfo以「缓存优先 → 查库 → 回写」提供统计。
三、面试常问知识点与难点
1. 存储引擎抽象接口设计(开闭原则)
类比:接口就像电源插座的标准。电器(业务)只认"插座长这样、提供 220V",至于背后是火电、水电还是电池(local/minio/oss),电器不在乎。新增一种电,只需做个新插头,不用改电器——这就是开闭原则(对扩展开放、对修改关闭)。
flowchart LR
Biz["UserUsecase / FileUsecase"] -->|依赖| Iface["biz.Storage 接口"]
Iface -.实现.-> Local["localStorage"]
Iface -.实现.-> Minio["minioStorage"]
Iface -.实现.-> Oss["ossStorage stub"]
Note["编译期断言
var _ biz.Storage = (*localStorage)(nil)
保证不漏方法"]代码里 var _ biz.Storage = (*localStorage)(nil) 就是给接口上的一道"保险丝"——只要 localStorage 漏实现任一方法,编译直接失败,把"忘写方法"挡在运行前。biz.Storage 接口把「业务要做什么」与「底层怎么做」解耦。新增存储引擎(如 OSS、COS、七牛云)时,只需新增一个实现并注册到工厂,业务代码零改动,符合开闭原则(对扩展开放,对修改关闭)。
2. 本地文件系统按日期分目录的意义
单目录文件数过多会显著降低文件系统性能:ext4 单目录超过 10 万文件后 readdir 变慢,inode 查找开销上升;备份、迁移、删除也更困难。按 年/月/日 分目录后,单目录文件数受限于每日上传量,且天然支持按时间归档与清理。
3. 零拷贝 io.Copy 优化
io.Copy 内部维护 32KB 缓冲区,从 Reader 读一块写一块,避免一次性把整个文件读入内存。对于大文件上传,这种方式内存占用恒定(O(1))。CompleteMultipartUpload 中更进一步使用 io.CopyBuffer 复用 1MB 缓冲区,避免每次合并分片都分配新内存,减少 GC 压力。
4. MinIO 与 S3 协议兼容性
MinIO 实现了 AWS S3 协议,minio-go SDK 同样可以连接 AWS S3、阿里云 OSS(S3 兼容模式)、腾讯云 COS 等。这意味着 minioStorage 实现天然具备多云迁移能力,只需更换 endpoint / accessKey / secretKey 即可切换云厂商。
5. 孤儿分片清理机制
分片上传是长流程,任何一步崩溃都会留下孤儿分片。CleanupOrphanChunks 通过 ModTime 判断过期,扫描 chunks/ 目录批量删除。本地实现用 os.ReadDir + os.RemoveAll,MinIO 实现用 ListObjects + RemoveObjects 批量接口,避免逐个删除的网络开销。
6. 存储路径生成策略
路径生成需保证:唯一性(避免覆盖)、可追溯(能从路径看出上传时间)、分散性(避免热点)。本项目用 日期目录/uuid.ext 或 日期前缀/uuid.merged:日期保证分散与可追溯,uuid 保证唯一。MinIO 中 objectKey 同样遵循此规则。
7. 分片合并的两种范式
类比:把几段录音拼成一首歌。本地方案是你把每段录音下载到自己电脑,再依次拼接(占你电脑内存和带宽);MinIO 方案是云服务商在云端直接把几段拼好,你只拿最终结果——服务端合并,不占应用服务器资源。
flowchart TB
subgraph 本地存储
P1[分片1] --> B1[io.CopyBuffer 顺序写]
P2[分片2] --> B1
P3[分片N] --> B1
B1 --> F1[最终文件 + Sync 落盘]
end
subgraph MinIO
Q1[分片对象1] --> C[ComposeObject 服务端合并]
Q2[分片对象2] --> C
Q3[分片对象N] --> C
C --> F2[合并后对象
不经过应用服务器]
end- 本地存储:按
PartNumber排序后,逐个os.Open分片,用io.CopyBuffer顺序写入最终文件,最后dest.Sync()强制落盘。 - MinIO 存储:使用
ComposeObject在对象存储服务端合并多个对象,不经过应用服务器,节省带宽与内存,是真正的「服务端合并」。
8. 配额原子更新与"并发超卖"这道坎
类比:统计已用空间像「公共计数器」,十个人同时加减,若各自先抄数再改,最后一定有人白改。正确做法是让数据库在一条 SQL 里「当场读当场改并上锁」,谁也插不了队。但还要注意:计数器只能往小改不能改负(删除防超卖),也不能往大改越过天花板(上传防超配额)——后者正是本模块教学版漏掉的地方。
flowchart LR
U[上传+ / 删除-] --> Q[UPDATE used_storage = used_storage + ?
WHERE id=? AND 条件]
Q --> R{RowsAffected}
R -->|1| OK[更新成功 无漂移]
R -->|0| X{用户存在?}
X -->|否| N[ErrUserNotFound]
X -->|是| S[ErrStorageInsufficient 空间不足]文件上传/删除会修改 used_storage。若先读后写,并发场景下会丢失更新。已落地的做法:
- 扣减(
delta<0)用WHERE id = ? AND used_storage >= ?(即used_storage >= -delta)做原子防超卖——对应 01-user.md 5.8 已讲。 - 增加(
delta>0,即上传)教学版当前只用WHERE id = ?,没有上界条件,存在并发超卖窗口(详见 5.6.2)。商用版必须补上WHERE id = ? AND used_storage + ? <= total_storage,让"配额天花板"在数据库这一层被强制守住。
9. 分层架构与依赖倒置(回顾)
biz 层定义 UserRepo、Cache、Locker 接口,data 层提供 GORM / Redis 实现,由 Wire 注入。biz 不直接依赖 gorm 或具体 cache 包,便于替换底层存储(如换 Postgres、换本地内存缓存),也方便单元测试时 mock。配额记账位于 biz 的 UserUsecase / Storage 辅助方法,真正落库在 data.userRepo,缓存与锁通过接口注入。
四、亿级流量优化思路
类比:小饭馆(单机本地存储)顾客多了就排队,想扩容得把一部分生意"外包"给中央厨房(对象存储 + CDN),自己只管接单和记账(元数据管理),压力瞬间小一圈。记账(配额)这条线同理:缓存挡读、DB 条件挡错、分布式锁挡冲突,三层各司其职。
flowchart LR
subgraph 传统模式
C1[客户端] --> S1[应用服务器 转发]
S1 --> O1[OSS]
end
subgraph 优化模式
C2[客户端] -->|预签名URL直传| O2[OSS 对象存储]
O2 --> CDN[CDN 边缘缓存]
C2 --> S2[应用服务器
只管元数据 + 配额]
end针对存储模块在高并发、大文件、海量用户场景下的优化方向:
对象存储 + CDN 加速:将热点文件从本地磁盘迁移到 OSS / MinIO,前置 CDN 边缘节点。用户下载命中 CDN 边缘缓存,回源流量大幅下降,延迟从跨地域降到同城。
GetFileURL返回的预签名 URL 可直接作为 CDN 回源地址。分片上传直传 OSS:传统模式是「客户端 → 应用服务器 → OSS」,应用服务器是带宽瓶颈。改造为「客户端拿到 STS 临时凭证 / 预签名 URL → 直传 OSS」,应用服务器只做元数据管理,带宽与 CPU 压力转移给云厂商。
存储路径分散避免单目录文件过多:按
年/月/日甚至年/月/日/小时分目录,单目录文件数控制在万级。对象存储虽无目录概念,但相同前缀的对象会落到同一分片(partition),前缀分散可避免热点分片。异步删除:用户删除文件时,先标记数据库
deleted_at,把对象删除任务投递到 Kafka / 消息队列,由消费者批量删除,避免同步删除拖慢接口响应。MinIO 已用go s.cleanupChunks(context.Background(), uploadID)起独立 goroutine 异步清理分片。多级存储冷热分离:根据文件访问频率自动分层。热数据放 SSD / OSS 标准存储,温数据放 OSS 低频存储,冷数据放 OSS 归档存储。通过
Lifecycle规则自动迁移,成本可下降一个数量级。分片合并服务端化:本地存储合并分片需经过应用服务器内存与磁盘,MinIO 用
ComposeObject在服务端合并,彻底消除应用服务器带宽瓶颈。亿级流量下应优先选择支持服务端合并的对象存储。断点续传与秒传:客户端记录已上传分片列表,断网后从断点继续;服务端通过文件 MD5 / 分片 Hash 判断是否已存在相同文件,命中则直接「秒传」建引用,无需重复上传。
配额记账的商用强化(新增,对应本模块审查重点):
- DB 上界兜底:增加的
UPDATE必须带used_storage + ? <= total_storage,即使应用层锁缺失也不超卖。 - 多级缓存防读崩:
GetStorageInfo引入 L1 本地缓存(如ristretto)→ L2 Redis → DB,Redis 缓存 key 加命名空间前缀、TTL 加 jitter,缓存空值防穿透、singleflight 防击穿。 - 锁只挡冲突不挡错误:分布式锁串行化同一用户,但正确性的唯一裁判是 DB 条件更新。
- 可观测:暴露用量率指标,超 80% 预警、超 95% 临界提示;
RecalibrateStorage漂移超阈值告警。 - 大用户与分库分表:校准用游标分页、避免
ListAllUserIDs全量载入;分库后用户行仍在单分片,行锁原子性不变,跨用户汇总走独立汇总表。
- DB 上界兜底:增加的
五、详细实现流程与代码解析
5.1 Storage 接口定义(抽象层设计)
实现思路
Storage 接口位于 internal/biz/storage.go,是业务层与存储层的契约。它定义了普通文件操作(上传 / 下载 / 删除 / 获取 URL)和分片上传全生命周期(初始化 / 上传分片 / 完成 / 取消)以及孤儿分片清理共 9 个方法。接口还配套了 PartInfo 结构体描述分片元数据,以及一组通用错误(ErrStorageNotImplemented、ErrStorageFileNotFound 等)作为各实现层的统一错误语义。
这种抽象的好处是:业务层 UserUsecase、FileUsecase 等只持有 Storage 接口字段,通过依赖注入可在本地磁盘、MinIO、OSS 之间无缝切换,符合依赖倒置原则。
关键代码
// 文件: internal/biz/storage.go
// Storage 定义了文件存储操作的接口
// biz 层仅依赖此接口——不依赖任何具体的存储实现
type Storage interface {
// Upload 从读取器存储文件并返回存储路径
// 返回的 storagePath 用于后续操作中标识文件
Upload(ctx context.Context, filePath string, reader io.Reader) (storagePath string, err error)
// Download 返回给定存储路径的文件的 ReadCloser
// 调用方使用完毕后必须 Close 返回的 ReadCloser,避免文件描述符泄漏
Download(ctx context.Context, storagePath string) (io.ReadCloser, error)
// Delete 从存储中永久删除文件
// 文件不存在时实现层应返回 nil,保证幂等性
Delete(ctx context.Context, storagePath string) error
// InitMultipartUpload 初始化分片上传会话
// 使用调用方传入的 uploadID 作为分片目录标识,确保与 biz 层会话 ID 一致
// 这样设计避免了 storage 层自己生成 ID 与 biz 层不一致导致的孤儿空目录
InitMultipartUpload(ctx context.Context, uploadID, fileName string) error
// UploadPart 上传分片上传的单个分片
// 返回已上传分片的元数据(PartNumber / Hash / Size)
UploadPart(ctx context.Context, uploadID string, partNumber int, reader io.Reader) (PartInfo, error)
// CompleteMultipartUpload 完成分片上传并合并所有分片
// 返回合并后文件的最终存储路径
CompleteMultipartUpload(ctx context.Context, uploadID string, parts []PartInfo) (storagePath string, err error)
// AbortMultipartUpload 取消分片上传并清理临时分片
AbortMultipartUpload(ctx context.Context, uploadID string) error
// CleanupOrphanChunks 清理孤儿分片目录
// 扫描 chunks 目录,删除修改时间早于 maxAge 的子目录
// 返回被清理的 uploadID 列表,便于上层记录日志或回调
CleanupOrphanChunks(ctx context.Context, maxAge time.Duration) ([]string, error)
// GetFileURL 返回可用于文件预览或分享的访问 URL
// 本地存储返回本地路径;OSS 返回签名 URL
GetFileURL(ctx context.Context, storagePath string, expire time.Duration) (string, error)
}
// PartInfo 保存已上传分片的元数据
// CompleteMultipartUpload 时会收到 []PartInfo 用于按序合并
type PartInfo struct {
PartNumber int // 分片编号,从 1 开始,决定合并顺序
Hash string // 分片哈希 / ETag,用于校验完整性
Size int64 // 分片字节数
}
// 常见的存储错误
// 所有实现层共用这套错误,便于 biz 层用 errors.Is 精确判断
var (
ErrStorageNotImplemented = errors.New("storage: not implemented") // OSS stub 默认返回此错误
ErrStorageFileNotFound = errors.New("storage: file not found") // 下载 / 删除时文件不存在
ErrUploadNotFound = errors.New("storage: upload session not found")
ErrPartAlreadyUploaded = errors.New("storage: part already uploaded") // 同一分片重复上传
)
工厂选择代码(internal/data/storage/storage.go)根据配置 Driver 字段决定实例化哪个实现:
// 文件: internal/data/storage/storage.go
// ProviderSet 是存储引擎的 Wire 提供集合。
// Wire 在编译期生成依赖注入代码时,会从此集合中找到 NewStorage 并注入。
var ProviderSet = wire.NewSet(NewStorage)
// NewStorage 根据配置中的 driver 字段选择对应的存储实现。
// 支持 "local" 和 "minio" 两种驱动。
// 后续扩展 OSS 时,只需在此 switch 增加 case "oss": return NewOssStorage(c.Oss)
func NewStorage(c *conf.Storage) (biz.Storage, error) {
if c == nil {
return nil, fmt.Errorf("storage: 存储配置为空")
}
switch c.Driver {
case "minio":
return NewMinioStorage(c.Minio)
case "local", "":
return NewLocalStorage(c.Local), nil
default:
return nil, fmt.Errorf("storage: 不支持的存储驱动: %s", c.Driver)
}
}
5.2 本地文件存储实现(上传 / 下载 / 删除 / 路径生成)
实现思路
localStorage 用 Go 标准库 os / io 直接操作文件系统,适合单机部署与开发测试。核心设计:
- 构造函数:
NewLocalStorage从conf.Storage_Local读取upload_dir和max_file_size,默认./data/uploads与 100MB,启动时os.MkdirAll确保目录存在。 - 路径生成:
datePath()用time.Now()拼出uploadDir/2026/08/04/,每次上传都落到当天目录,天然分散。 - 上传:
os.Create创建目标文件,io.Copy流式写入,避免大文件占满内存。 - 下载:
os.Open返回*os.File(实现io.ReadCloser),文件不存在时返回biz.ErrStorageFileNotFound。 - 删除:
os.Remove,文件不存在视为成功(幂等)。 - URL:本地无 HTTP 服务,直接返回
file://绝对路径便于本地预览。
关键代码
// 文件: internal/data/storage/local.go
// 编译期确保 localStorage 实现了 biz.Storage 接口。
// 若接口新增方法而 localStorage 未实现,编译阶段就会报错,避免运行时才发现遗漏。
var _ biz.Storage = (*localStorage)(nil)
// localStorage 使用本地文件系统实现 Storage 接口。
type localStorage struct {
uploadDir string // 上传根目录,例如 ./data/uploads
maxFileSize int64 // 单文件最大字节数
}
// NewLocalStorage 创建一个新的本地文件系统存储处理器。
// c 为空时使用默认值 ./data/uploads 与 100MB
func NewLocalStorage(c *conf.Storage_Local) biz.Storage {
dir := "./data/uploads"
var maxSize int64 = 100 * 1024 * 1024 // 默认 100MB
if c != nil {
if c.UploadDir != "" {
dir = c.UploadDir
}
if c.MaxFileSize != "" {
fmt.Sscanf(c.MaxFileSize, "%d", &maxSize) // 从字符串解析为 int64
}
}
// 确保上传目录存在,0755 表示 owner 可读写执行,其他人可读可执行
os.MkdirAll(dir, 0755)
return &localStorage{
uploadDir: dir,
maxFileSize: maxSize,
}
}
// datePath 返回按日期分割的路径:uploadDir/2026/08/04/
// 使用 filepath.Join 保证跨平台路径分隔符正确(Windows 用 \,Linux/Mac 用 /)
func (s *localStorage) datePath() string {
now := time.Now()
return filepath.Join(s.uploadDir,
fmt.Sprintf("%d", now.Year()), // 2026
fmt.Sprintf("%02d", now.Month()), // 08
fmt.Sprintf("%02d", now.Day()), // 04
)
}
// Upload 从读取器存储文件到本地文件系统。
// filePath 通常为 "uuid.extension",作为当天目录内的文件名
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)
}
// 使用 filePath 作为日期目录内的相对文件名
// filePath 通常为 "uuid.extension"
destPath := filepath.Join(dir, filePath)
// os.Create 若文件已存在会截断,否则新建
f, err := os.Create(destPath)
if err != nil {
return "", fmt.Errorf("storage: failed to create file: %w", err)
}
defer f.Close()
// io.Copy 内部维护 32KB 缓冲区,从 reader 读一块写一块
// 整个过程不需要把文件全部加载到内存,适合大文件
if _, err := io.Copy(f, reader); err != nil {
return "", fmt.Errorf("storage: failed to write file: %w", err)
}
return destPath, nil
}
// Download 返回给定存储路径的文件的 ReadCloser。
// 调用方必须 Close 返回值,否则文件描述符会泄漏
func (s *localStorage) Download(ctx context.Context, storagePath string) (io.ReadCloser, error) {
f, err := os.Open(storagePath)
if err != nil {
if os.IsNotExist(err) {
// 文件不存在时返回业务错误,便于上层用 errors.Is 判断
return nil, biz.ErrStorageFileNotFound
}
return nil, fmt.Errorf("storage: failed to open file: %w", err)
}
return f, nil
}
// Delete 从文件系统中删除文件。
// 文件不存在直接返回 nil,保证删除接口幂等
func (s *localStorage) Delete(ctx context.Context, storagePath string) error {
if err := os.Remove(storagePath); err != nil {
if os.IsNotExist(err) {
return nil // 已删除,不算错误
}
return fmt.Errorf("storage: failed to delete file: %w", err)
}
return nil
}
// GetFileURL 返回包装为 file:// URL 的本地文件路径。
// 本地存储没有 HTTP 服务,因此用 file:// 协议让浏览器或客户端直接访问本地路径
func (s *localStorage) GetFileURL(ctx context.Context, storagePath string, expire time.Duration) (string, error) {
if storagePath == "" {
return "", biz.ErrFileNotFound
}
// 转为绝对路径,避免相对路径在不同工作目录下失效
absPath, err := filepath.Abs(storagePath)
if err != nil {
return "", err
}
return "file://" + absPath, nil
}
// ChunkPath 返回特定分片文件的路径。
// 此方法非接口方法,提供给需要直接读取分片的内部逻辑使用
func (s *localStorage) ChunkPath(uploadID string, partNumber int) string {
// %05d 表示左补零至 5 位,例如 part_00001、part_00012
// 这样字典序与数字序一致,便于按文件名排序
return filepath.Join(s.uploadDir, "chunks", uploadID, fmt.Sprintf("part_%05d", partNumber))
}
5.3 MinIO 存储实现
类比:MinIO 不像本地硬盘有"文件夹",它更像一个有固定编号规则的超大仓库——
2026/08/04/uuid.ext是一个完整货位编号,中间的/只是给人看的分类标签,仓库本身不分层级。所以分片合并时它能"在仓库里直接拼",不用把货搬回你的店。
flowchart TB
A[Upload] -->|PutObject| K1["objectKey: 日期/uuid.ext"]
B[UploadPart] -->|每个分片独立对象| K2["chunks/{uploadID}/part_00001"]
C[CompleteMultipartUpload] -->|ComposeObject 服务端合并| K3["日期/{uploadID}.merged"]
C -->|go cleanupChunks| D[异步删除分片对象]实现思路
minioStorage 通过 minio-go/v7 SDK 与 MinIO / S3 兼容对象存储交互。与本地存储相比:
- 对象 key 代替文件路径:MinIO 没有「目录」概念,
2026/08/04/uuid.ext是一个完整对象 key,前缀/仅用于虚拟分组。 - 构造函数:
NewMinioStorage创建 client,启动时检查bucket是否存在,不存在则自动创建,降低部署门槛。 - 上传 / 下载 / 删除:分别用
PutObject/GetObject/RemoveObject。下载时先obj.Stat()验证对象存在,避免返回一个读取出错才暴露的流。 - 分片上传:每个分片作为独立对象
chunks/{uploadID}/part_00001上传,合并时用ComposeObject在服务端把多个对象拼成一个,完全不经过应用服务器内存。 - 预签名 URL:
PresignedGetObject生成带过期时间的下载链接,前端可直接用该 URL 访问 MinIO,无需经过应用服务器转发。 - 异步清理:
CompleteMultipartUpload合并成功后,用go s.cleanupChunks(context.Background(), uploadID)起独立 goroutine 删除分片对象,不阻塞主流程返回。注意 context 用context.Background()而非请求 ctx,避免请求结束后清理被取消。
关键代码
// 文件: internal/data/storage/minio.go
// 编译期确保 minioStorage 实现了 biz.Storage 接口。
var _ biz.Storage = (*minioStorage)(nil)
// minioStorage 使用 MinIO 对象存储实现 Storage 接口。
type minioStorage struct {
client *minio.Client // MinIO / S3 客户端
bucket string // bucket 名称,例如 cloud-disk
location string // region,例如 us-east-1
}
// NewMinioStorage 创建一个新的 MinIO 存储处理器。
// 启动时自动检查并创建 bucket,避免部署时手动准备
func NewMinioStorage(c *conf.Storage_Minio) (biz.Storage, error) {
if c == nil {
return nil, fmt.Errorf("storage: minio 配置为空")
}
useSSL := c.UseSsl
// credentials.NewStaticV4 使用静态 AK/SK 签名
client, err := minio.New(c.Endpoint, &minio.Options{
Creds: credentials.NewStaticV4(c.AccessKey, c.SecretKey, ""),
Secure: useSSL,
Region: c.Region,
})
if err != nil {
return nil, fmt.Errorf("storage: 创建 minio 客户端失败: %w", err)
}
bucket := c.BucketName
if bucket == "" {
bucket = "cloud-disk" // 默认 bucket
}
location := c.Region
if location == "" {
location = "us-east-1" // S3 默认 region
}
// 启动时确保 bucket 存在,避免运行时才发现配置错误
ctx := context.Background()
exists, err := client.BucketExists(ctx, bucket)
if err != nil {
return nil, fmt.Errorf("storage: 检查 minio bucket 失败: %w", err)
}
if !exists {
if err := client.MakeBucket(ctx, bucket, minio.MakeBucketOptions{Region: location}); err != nil {
return nil, fmt.Errorf("storage: 创建 minio bucket 失败: %w", err)
}
slog.Info("minio bucket 创建成功", "bucket", bucket)
}
slog.Info("minio 存储初始化成功", "endpoint", c.Endpoint, "bucket", bucket)
return &minioStorage{
client: client,
bucket: bucket,
location: location,
}, nil
}
// dateKey 返回按日期分割的对象前缀:2026/08/04/
// 注意 MinIO 用的是 / 分隔的字符串,而非 filepath.Join(避免 Windows 平台 \ 污染)
func (s *minioStorage) dateKey() string {
now := time.Now()
return fmt.Sprintf("%d/%02d/%02d", now.Year(), int(now.Month()), now.Day())
}
// Upload 从读取器存储文件到 MinIO。
// objectKey 格式:2026/08/04/uuid.ext
// size 传 -1 表示未知大小,MinIO SDK 会分块上传
func (s *minioStorage) Upload(ctx context.Context, filePath string, reader io.Reader) (string, error) {
objectKey := fmt.Sprintf("%s/%s", s.dateKey(), filePath)
_, err := s.client.PutObject(ctx, s.bucket, objectKey, reader, -1, minio.PutObjectOptions{
ContentType: "application/octet-stream", // 二进制流默认类型
})
if err != nil {
return "", fmt.Errorf("storage: minio 上传文件失败: %w", err)
}
return objectKey, nil
}
// Download 返回给定存储路径的文件的 ReadCloser。
// 注意:GetObject 返回的 obj 在实际读取前不会报错,必须调用 Stat() 验证
func (s *minioStorage) Download(ctx context.Context, storagePath string) (io.ReadCloser, error) {
obj, err := s.client.GetObject(ctx, s.bucket, storagePath, minio.GetObjectOptions{})
if err != nil {
return nil, fmt.Errorf("storage: minio 获取对象失败: %w", err)
}
// 尝试读取以验证对象是否存在
if _, err := obj.Stat(); err != nil {
obj.Close()
resp := minio.ToErrorResponse(err)
if resp.Code == "NoSuchKey" {
return nil, biz.ErrStorageFileNotFound // 转为业务错误
}
return nil, fmt.Errorf("storage: minio 获取对象信息失败: %w", err)
}
return obj, nil
}
// Delete 从 MinIO 中删除文件。
// NoSuchKey 视为成功,保证幂等
func (s *minioStorage) Delete(ctx context.Context, storagePath string) error {
err := s.client.RemoveObject(ctx, s.bucket, storagePath, minio.RemoveObjectOptions{})
if err != nil {
resp := minio.ToErrorResponse(err)
if resp.Code == "NoSuchKey" {
return nil
}
return fmt.Errorf("storage: minio 删除对象失败: %w", err)
}
return nil
}
// InitMultipartUpload 初始化分片上传会话。
// MinIO 方式下使用对象前缀模拟分片目录,无需额外初始化。
// 放入一个标记对象 .init 以标识会话存在,便于排查
func (s *minioStorage) InitMultipartUpload(ctx context.Context, uploadID, fileName string) error {
markerKey := fmt.Sprintf("chunks/%s/.init", uploadID)
// 写入 uploadID 字符串作为标记内容,便于后续排查
_, err := s.client.PutObject(ctx, s.bucket, markerKey, strings.NewReader(uploadID), int64(len(uploadID)), minio.PutObjectOptions{})
if err != nil {
return fmt.Errorf("storage: minio 初始化分片上传失败: %w", err)
}
return nil
}
// UploadPart 上传分片上传的单个分片。
// 每个分片是一个独立对象:chunks/{uploadID}/part_00001
func (s *minioStorage) UploadPart(ctx context.Context, uploadID string, partNumber int, reader io.Reader) (biz.PartInfo, error) {
partKey := fmt.Sprintf("chunks/%s/part_%05d", uploadID, partNumber)
// 检查分片是否已存在(断点续传时同一分片可能重复上传)
_, err := s.client.StatObject(ctx, s.bucket, partKey, minio.StatObjectOptions{})
if err == nil {
return biz.PartInfo{}, biz.ErrPartAlreadyUploaded // 已上传,拒绝重复写
}
// info.ETag 是 MinIO 返回的分片哈希,合并时可作为校验依据
info, err := s.client.PutObject(ctx, s.bucket, partKey, reader, -1, minio.PutObjectOptions{
ContentType: "application/octet-stream",
})
if err != nil {
return biz.PartInfo{}, fmt.Errorf("storage: minio 上传分片失败: %w", err)
}
return biz.PartInfo{
PartNumber: partNumber,
Hash: info.ETag, // S3 协议中 ETag 通常即分片 MD5
Size: info.Size,
}, nil
}
// CompleteMultipartUpload 将所有已上传的分片合并为最终文件。
// 使用 MinIO 的 ComposeObject API 将多个分片对象在服务端合并为一个,
// 整个过程不经过应用服务器内存,是真正的「服务端合并」
func (s *minioStorage) CompleteMultipartUpload(ctx context.Context, uploadID string, parts []biz.PartInfo) (string, error) {
if len(parts) == 0 {
return "", fmt.Errorf("storage: 没有分片可合并")
}
// 按分片编号排序,保证合并后的字节顺序正确
sort.Slice(parts, func(i, j int) bool {
return parts[i].PartNumber < parts[j].PartNumber
})
// 构建最终对象路径:2026/08/04/{uploadID}.merged
finalKey := fmt.Sprintf("%s/%s.merged", s.dateKey(), uploadID)
// 构建 Compose 源:每个分片对象作为一个 CopySrcOptions
sources := make([]minio.CopySrcOptions, len(parts))
for i, part := range parts {
sources[i] = minio.CopySrcOptions{
Bucket: s.bucket,
Object: fmt.Sprintf("chunks/%s/part_%05d", uploadID, part.PartNumber),
}
}
// ComposeObject 在 MinIO 服务端把多个对象顺序拼接成一个新对象
dst := minio.CopyDestOptions{
Bucket: s.bucket,
Object: finalKey,
}
_, err := s.client.ComposeObject(ctx, dst, sources...)
if err != nil {
return "", fmt.Errorf("storage: minio 合并分片失败: %w", err)
}
// 清理分片对象 — 使用独立后台 context 避免请求取消后清理中断
// 例如客户端取消请求后请求 ctx 会被 cancel,若用请求 ctx 清理会半途中断
go s.cleanupChunks(context.Background(), uploadID)
return finalKey, nil
}
// AbortMultipartUpload 取消分片上传并清理临时分片。
func (s *minioStorage) AbortMultipartUpload(ctx context.Context, uploadID string) error {
return s.cleanupChunks(ctx, uploadID)
}
// cleanupChunks 清理指定 uploadID 的所有分片对象。
// 使用 RemoveObjects 批量接口,避免逐个删除的网络往返开销
func (s *minioStorage) cleanupChunks(ctx context.Context, uploadID string) error {
prefix := fmt.Sprintf("chunks/%s/", uploadID)
// ListObjects 返回一个 channel,流式列举,适合海量对象
objectsCh := s.client.ListObjects(ctx, s.bucket, minio.ListObjectsOptions{
Prefix: prefix,
Recursive: true, // 递归列举子「目录」
})
// RemoveObjects 接收对象 channel,批量删除并返回错误 channel
for err := range s.client.RemoveObjects(ctx, s.bucket, objectsCh, minio.RemoveObjectsOptions{}) {
if err.Err != nil {
slog.Warn("清理 MinIO 分片对象失败", "uploadID", uploadID, "err", err.Err)
}
}
return nil
}
// GetFileURL 返回 MinIO 预签名 URL,用于文件预览或分享。
// 预签名 URL 内嵌签名与过期时间,持有者无需 AK/SK 即可下载
func (s *minioStorage) GetFileURL(ctx context.Context, storagePath string, expire time.Duration) (string, error) {
if storagePath == "" {
return "", biz.ErrFileNotFound
}
// PresignedGetObject 生成 GET 方法的预签名 URL
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
}
5.4 阿里云 OSS 扩展规范
实现思路
oss_stub.go 是 OSS 实现的「规范文档 + 桩代码」。它定义了 ossStorage 结构体与所有接口方法的签名,但每个方法都返回 biz.ErrStorageNotImplemented。文件头部的注释列出了:
- 依赖:
github.com/aliyun/aliyun-oss-go-sdk/v3 - 配置字段:
Endpoint/AccessKeyId/AccessKeySecret/BucketName/Region - 每个方法对应的 OSS SDK 调用:例如
Upload对应bucket.PutObject,CompleteMultipartUpload对应bucket.CompleteMultipartUpload,GetFileURL对应bucket.SignURL。
这种「先写规范后填实现」的模式让团队可以并行工作:业务层立刻能基于 biz.Storage 接口开发,OSS 实现可以后续按 stub 注释补全,无需阻塞主流程。NewStorage 工厂只需增加 case "oss": return NewOssStorage(c.Oss) 即可启用。
关键代码
// 文件: internal/data/storage/oss_stub.go
package storage
import (
"context"
"io"
"time"
"cloud-disk/internal/biz"
)
// Ensure ossStorage implements biz.Storage at compile time.
// 编译期断言:确保 ossStorage 实现了 biz.Storage 接口
// 即便所有方法都返回 ErrStorageNotImplemented,也必须实现全部方法签名
var _ biz.Storage = (*ossStorage)(nil)
// ossStorage provides a stub implementation of the Storage interface for Aliyun OSS.
// TODO(oss): Replace stub with real implementation using github.com/aliyun/aliyun-oss-go-sdk
//
// Required dependencies (add to go.mod when implementing):
// github.com/aliyun/aliyun-oss-go-sdk/v3
//
// Configuration fields (from conf.Storage_Oss):
// - Endpoint: OSS endpoint URL (e.g., "oss-cn-hangzhou.aliyuncs.com")
// - AccessKeyId: Aliyun RAM access key ID
// - AccessKeySecret: Aliyun RAM access key secret
// - BucketName: OSS bucket name
// - Region: OSS region (e.g., "cn-hangzhou")
//
// Key implementation notes:
// - Upload: Use bucket.PutObject() with io.Reader
// - Download: Use bucket.GetObject() returning io.ReadCloser
// - Delete: Use bucket.DeleteObject()
// - InitMultipartUpload: Use bucket.InitiateMultipartUpload()
// - UploadPart: Use bucket.UploadPart() with upload ID and part number
// - CompleteMultipartUpload: Use bucket.CompleteMultipartUpload() with all ETags
// - AbortMultipartUpload: Use bucket.AbortMultipartUpload()
// - GetFileURL: Use bucket.SignURL() with the given expiration duration
type ossStorage struct {
// TODO(oss): Add fields:
// client *oss.Client
// bucket *oss.Bucket
// bucketName string
}
// NewOssStorage creates a new OSS storage handler stub.
// TODO(oss): Initialize the OSS client with credentials from conf.
// 后续实现时签名应改为 NewOssStorage(c *conf.Storage_Oss) (biz.Storage, error)
func NewOssStorage() biz.Storage {
return &ossStorage{}
}
// Upload — TODO: bucket.PutObject(filePath, reader)
func (s *ossStorage) Upload(ctx context.Context, filePath string, reader io.Reader) (string, error) {
return "", biz.ErrStorageNotImplemented
}
// Download — TODO: return bucket.GetObject(storagePath)
func (s *ossStorage) Download(ctx context.Context, storagePath string) (io.ReadCloser, error) {
return nil, biz.ErrStorageNotImplemented
}
// Delete — TODO: bucket.DeleteObject(storagePath)
func (s *ossStorage) Delete(ctx context.Context, storagePath string) error {
return biz.ErrStorageNotImplemented
}
// InitMultipartUpload — TODO: result, err := bucket.InitiateMultipartUpload(fileName)
// 使用 biz 层传入的 uploadID 作为本地会话标识
func (s *ossStorage) InitMultipartUpload(ctx context.Context, uploadID, fileName string) error {
return biz.ErrStorageNotImplemented
}
// UploadPart — TODO: part, err := bucket.UploadPart(initResult, reader, partSize, partNumber)
// return biz.PartInfo{PartNumber: partNumber, Hash: part.ETag, Size: partSize}, nil
func (s *ossStorage) UploadPart(ctx context.Context, uploadID string, partNumber int, reader io.Reader) (biz.PartInfo, error) {
return biz.PartInfo{}, biz.ErrStorageNotImplemented
}
// CompleteMultipartUpload — TODO: Build oss.Parts from parts,
// call bucket.CompleteMultipartUpload(initResult, ossParts), return storagePath, nil
func (s *ossStorage) CompleteMultipartUpload(ctx context.Context, uploadID string, parts []biz.PartInfo) (string, error) {
return "", biz.ErrStorageNotImplemented
}
// AbortMultipartUpload — TODO: bucket.AbortMultipartUpload(initResult)
func (s *ossStorage) AbortMultipartUpload(ctx context.Context, uploadID string) error {
return biz.ErrStorageNotImplemented
}
// CleanupOrphanChunks — TODO: 列举并清理 OSS 上的孤儿分片对象(前缀 chunks/),
// 可通过 ListObjectsV2 + DeleteObjects 实现。
func (s *ossStorage) CleanupOrphanChunks(ctx context.Context, maxAge time.Duration) ([]string, error) {
return nil, biz.ErrStorageNotImplemented
}
// GetFileURL — TODO: signedURL, err := bucket.SignURL(storagePath, oss.HTTPGet, expire)
// return signedURL, nil
func (s *ossStorage) GetFileURL(ctx context.Context, storagePath string, expire time.Duration) (string, error) {
return "", biz.ErrStorageNotImplemented
}
配置定义(internal/conf/conf.proto)已为 OSS 预留字段:
// 文件: internal/conf/conf.proto
message Storage {
string driver = 1; // 存储驱动:local / minio / oss
message Local {
string upload_dir = 1;
string max_file_size = 2;
int32 concurrent_upload = 3;
}
message Oss { // 阿里云 OSS 配置
string endpoint = 1; // 例如 oss-cn-hangzhou.aliyuncs.com
string access_key_id = 2;
string access_key_secret = 3;
string bucket_name = 4;
string region = 5; // 例如 cn-hangzhou
}
message Minio {
string endpoint = 1;
string access_key = 2;
string secret_key = 3;
string bucket_name = 4;
string region = 5;
bool use_ssl = 6;
}
Local local = 2;
Oss oss = 3;
Minio minio = 4;
}
5.5 孤儿分片目录清理
类比:停车场里的"僵尸车位"——某辆车(分片上传会话)开到一半人跑了,车位一直占着。管理员定期巡查,超过 N 天没人动的车位就清掉,腾出空间。
flowchart TB
A[分片上传开始
InitMultipartUpload] --> B[多次 UploadPart]
B --> C{完成?}
C -->|Complete| D[合并并清理分片]
C -->|崩溃/断开/放弃| E[孤儿分片残留]
E --> F[CleanupOrphanChunks 定时扫描]
F --> G{ModTime > maxAge?}
G -->|是| H[删除分片目录/对象]
G -->|否| I[保留]实现思路
分片上传是长流程:客户端先 InitMultipartUpload,多次 UploadPart,最后 CompleteMultipartUpload。任何一步因网络中断、服务崩溃、客户端放弃都会留下未清理的分片,称为「孤儿分片」。CleanupOrphanChunks 的职责是定期扫描 chunks/ 目录,删除存活时间超过 maxAge 的分片,回收存储空间。
本地存储清理逻辑:
- 用
os.ReadDir列举uploadDir/chunks/下所有条目(比ioutil.ReadDir更高效,不一次性调用Stat)。 - 对每个目录条目调用
entry.Info()获取ModTime。 - 若
now.Sub(ModTime) > maxAge,用os.RemoveAll删除整个分片目录。 - 删除失败不中断整体流程,只打印日志,继续处理其他目录。
- 全程检查
ctx.Err(),支持调用方取消。 - 返回被清理的
uploadID列表,便于上层记录或回调。
MinIO 清理逻辑:
ListObjects流式列举chunks/前缀下所有对象。- 收集
LastModified早于now - maxAge的对象。 - 从对象 key
chunks/{uploadID}/part_xxx提取uploadID(用strings.SplitN切 3 段取第 2 段)。 - 用
RemoveObjects批量删除,避免逐个删除的网络开销。 - 返回去重后的
uploadID列表。
关键代码
本地存储的孤儿分片清理:
// 文件: internal/data/storage/local.go
// CleanupOrphanChunks 扫描 chunks 目录并删除修改时间早于 maxAge 的子目录。
// 这些残留目录通常由服务异常崩溃或客户端断开未调用取消接口导致。
// 返回被清理的 uploadID 列表。
func (s *localStorage) CleanupOrphanChunks(ctx context.Context, maxAge time.Duration) ([]string, error) {
chunksRoot := filepath.Join(s.uploadDir, "chunks")
// os.ReadDir 是 Go 1.16+ 引入的高效列举方法,
// 返回 []os.DirEntry 而非 []os.FileInfo,不会预读取每个条目的完整信息
entries, err := os.ReadDir(chunksRoot)
if err != nil {
if os.IsNotExist(err) {
// chunks 目录不存在,说明没有任何分片上传发生过
return nil, nil
}
return nil, fmt.Errorf("storage: failed to read chunks dir: %w", err)
}
now := time.Now()
// 预分配容量,减少 append 时的扩容拷贝
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()
dirPath := filepath.Join(chunksRoot, uploadID)
// DirEntry.Info() 按需获取条目信息,比传统 ReadDir 更省资源
info, err := entry.Info()
if err != nil {
// 无法获取信息,跳过(不删除以避免误操作)
continue
}
// 使用目录的 ModTime 判断;超过 maxAge 视为孤儿分片
// ModTime 是目录最后被修改的时间(如分片写入会更新它)
if now.Sub(info.ModTime()) > maxAge {
if err := os.RemoveAll(dirPath); 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 的孤儿分片清理:
// 文件: internal/data/storage/minio.go
// CleanupOrphanChunks 扫描 chunks/ 前缀下超过 maxAge 的分片对象并删除。
func (s *minioStorage) CleanupOrphanChunks(ctx context.Context, maxAge time.Duration) ([]string, error) {
// 列举所有 chunks/ 下的对象
objectCh := s.client.ListObjects(ctx, s.bucket, minio.ListObjectsOptions{
Prefix: "chunks/",
Recursive: true, // 递归列举所有层级
})
now := time.Now()
uploadIDs := make(map[string]bool) // 用 map 去重 uploadID
var objectsToDelete []minio.ObjectInfo // 收集待删除对象
for obj := range objectCh {
if ctx.Err() != nil {
return nil, ctx.Err() // 支持取消
}
if obj.Err != nil {
continue // 列举错误跳过单个对象
}
// LastModified 是对象在 MinIO 上的最后修改时间
if now.Sub(obj.LastModified) > maxAge {
objectsToDelete = append(objectsToDelete, obj)
// 从对象路径提取 uploadID: chunks/{uploadID}/part_xxx
// SplitN 切 3 段:["chunks", "{uploadID}", "part_xxx"]
parts := strings.SplitN(obj.Key, "/", 3)
if len(parts) >= 2 {
uploadIDs[parts[1]] = true
}
}
}
// 批量删除过期对象
// RemoveObjects 接收一个 channel,内部并发删除,效率远高于循环 RemoveObject
if len(objectsToDelete) > 0 {
delCh := make(chan minio.ObjectInfo, len(objectsToDelete))
for _, obj := range objectsToDelete {
delCh <- obj
}
close(delCh) // 必须关闭,RemoveObjects 才能知道输入结束
for err := range s.client.RemoveObjects(ctx, s.bucket, delCh, minio.RemoveObjectsOptions{}) {
if err.Err != nil {
slog.Warn("清理 MinIO 孤儿分片失败", "err", err.Err)
}
}
}
// 把 map 的 key 转为 slice 返回
result := make([]string, 0, len(uploadIDs))
for id := range uploadIDs {
result = append(result, id)
}
return result, nil
}
本地存储分片合并代码补充(CompleteMultipartUpload,体现合并完成后的清理与高性能写入):
// 文件: internal/data/storage/local.go
// CompleteMultipartUpload 将所有已上传的分片合并为最终文件。
func (s *localStorage) CompleteMultipartUpload(ctx context.Context, uploadID string, parts []biz.PartInfo) (string, error) {
chunkDir := filepath.Join(s.uploadDir, "chunks", uploadID)
dir := s.datePath()
if err := os.MkdirAll(dir, 0755); err != nil {
return "", fmt.Errorf("storage: failed to create final dir: %w", err)
}
// 最终文件名用 uploadID + .merged 后缀,避免与普通上传文件混淆
finalPath := filepath.Join(dir, uploadID+".merged")
// 按分片编号排序,确保合并后字节顺序正确
sort.Slice(parts, func(i, j int) bool {
return parts[i].PartNumber < parts[j].PartNumber
})
dest, err := os.Create(finalPath)
if err != nil {
return "", fmt.Errorf("storage: failed to create final file: %w", err)
}
defer dest.Close()
// 使用带缓冲的 Writer 减少 write 系统调用次数
// 缓冲区 1MB,对于大文件合并可显著减少系统调用
bufWriter := bufio.NewWriterSize(dest, 1024*1024)
defer bufWriter.Flush()
// 复用 io.Copy 的缓冲区,避免每次复制都分配新内存
// 1MB 缓冲区在内存占用与拷贝效率之间取得平衡
copyBuf := make([]byte, 1024*1024)
// 按顺序合并分片
for _, part := range parts {
partPath := filepath.Join(chunkDir, fmt.Sprintf("part_%05d", part.PartNumber))
src, err := os.Open(partPath)
if err != nil {
return "", fmt.Errorf("storage: failed to open part %d: %w", part.PartNumber, err)
}
// io.CopyBuffer 使用调用方提供的缓冲区,避免内部分配
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()
}
// 刷新缓冲区并确保数据落盘
// Flush 把 bufio 缓冲区的数据写到 os.File;Sync 调用 fsync 强制刷盘
if err := bufWriter.Flush(); err != nil {
return "", fmt.Errorf("storage: failed to flush merged file: %w", err)
}
if err := dest.Sync(); err != nil {
return "", fmt.Errorf("storage: failed to sync merged file: %w", err)
}
// 清理分片目录 — 合并成功后分片已无用处
os.RemoveAll(chunkDir)
return finalPath, nil
}
5.6 配额与存储统计(商用强化)
为什么单独成节:存储引擎负责"货放哪儿",配额负责"账记清楚"。商用在线服务里,配额错误会直接导致用户被超卖(别人占满你的空间)或统计漂移(用了却没扣、删了却没减),必须当成一等公民对待。下面所有「修正点」都是教学版相对商用标准的缺口。
5.6.1 数据模型与配额字段
model.User 用三个 int64 字段描述容量(internal/data/model/user.go):
// 文件: internal/data/model/user.go
type User struct {
// ...
StorageQuota int64 `gorm:"column:storage_quota;default:10737418240;not null"` // bytes, default 10GB
TotalStorage int64 `gorm:"column:total_storage;default:10737418240;not null"` // bytes, default 10GB
UsedStorage int64 `gorm:"column:used_storage;default:0;not null"`
// ...
}
- 有效配额 =
total_storage:GetStorageInfo把total_storage当作Total(天花板),used_storage当作已用,Available = total - used。 storage_quota是冗余列:它与total_storage默认值相同(都 10GB),但存储统计逻辑从未读取storage_quota。- 修正点:商用设计应二选一——要么删除
storage_quota,只保留total_storage作为唯一配额源;要么明确分工:storage_quota为"套餐上限策略值",total_storage为"实际授予容量"(升级套餐时同步两者),并在GetStorageInfo统一用total_storage计算余量。当前以total_storage为准。
- 修正点:商用设计应二选一——要么删除
used_storage的唯一维护方式:UpdateUsedStorageAtomic的原子增减(见 5.6.2),任何业务路径都不得用Updates(map)直接写它,否则并发下必漂移。
5.6.2 原子扣减 UpdateUsedStorageAtomic(并发超卖的核心修复)
代码位于 internal/data/user.go。教学版现状:
// 文件: internal/data/user.go(教学版,含缺陷)
func (r *userRepo) UpdateUsedStorageAtomic(ctx context.Context, userID uint64, delta int64) error {
if delta < 0 {
// 减少:带 WHERE used_storage >= ? 防超卖 —— 正确
result := r.db.WithContext(ctx).Model(&model.User{}).
Where("id = ? AND used_storage >= ?", userID, -delta).
UpdateColumn("used_storage", gorm.Expr("used_storage + ?", delta))
if result.Error != nil {
return result.Error
}
if result.RowsAffected == 0 {
var count int64
r.db.WithContext(ctx).Model(&model.User{}).Where("id = ?", userID).Count(&count)
if count == 0 {
return biz.ErrUserNotFound
}
return biz.ErrStorageInsufficient
}
return nil
}
// 增加(上传):仅 WHERE id = ?,无上界 —— 缺陷!
result := r.db.WithContext(ctx).Model(&model.User{}).
Where("id = ?", userID).
UpdateColumn("used_storage", gorm.Expr("used_storage + ?", delta))
if result.Error != nil {
return result.Error
}
if result.RowsAffected == 0 {
return biz.ErrUserNotFound // 仅判断用户不存在,漏了"配额不足"
}
return nil
}
缺陷分析(并发超卖):增加分支没有 used_storage + delta <= total_storage 的约束。即使 CheckStorageAvailable 预检通过,两个并发上传也可能都基于陈旧的缓存余量放行,再相继 used_storage += delta,最终 used_storage 超过 total_storage——即用户被超卖,别人还能继续塞。单靠 updateStorageWithLock 的分布式锁并不够(见 5.6.5)。
修正点:给增加分支补上 DB 上界条件,让配额天花板在数据库这一层被强制守住,即便应用锁缺失也不会超卖:
// 文件: internal/data/user.go(商用修正版)
func (r *userRepo) UpdateUsedStorageAtomic(ctx context.Context, userID uint64, delta int64) error {
if delta < 0 {
// 减少:used_storage >= -delta 防扣成负数(超卖反向)
result := r.db.WithContext(ctx).Model(&model.User{}).
Where("id = ? AND used_storage >= ?", userID, -delta).
UpdateColumn("used_storage", gorm.Expr("used_storage + ?", delta))
if result.Error != nil {
return result.Error
}
if result.RowsAffected == 0 {
return r.classifyZero(ctx, userID) // 区分"用户不存在"与"空间不足"
}
return nil
}
// 修正点:增加也必须带 DB 上界,防止并发超卖
// 条件等价于 used_storage + delta <= total_storage
result := r.db.WithContext(ctx).Model(&model.User{}).
Where("id = ? AND used_storage <= total_storage - ?", userID, delta).
UpdateColumn("used_storage", gorm.Expr("used_storage + ?", delta))
if result.Error != nil {
return result.Error
}
if result.RowsAffected == 0 {
return r.classifyZero(ctx, userID) // 同样要区分两种 0 行情况
}
return nil
}
// classifyZero 在 RowsAffected==0 时区分"用户不存在"与"配额不足"
func (r *userRepo) classifyZero(ctx context.Context, userID uint64) error {
var count int64
r.db.WithContext(ctx).Model(&model.User{}).Where("id = ?", userID).Count(&count)
if count == 0 {
return biz.ErrUserNotFound
}
return biz.ErrStorageInsufficient
}
flowchart TD
U[UpdateUsedStorageAtomic delta] --> D{delta 正负?}
D -->|delta < 0 删除| A1[UPDATE used_storage = used - ?
WHERE id=? AND used >= ?]
D -->|delta > 0 上传| A2[UPDATE used_storage = used + ?
WHERE id=? AND used <= total - ?]
A1 --> R{RowsAffected}
A2 --> R
R -->|1| OK[成功 无漂移]
R -->|0| C{用户存在?}
C -->|否| N[ErrUserNotFound]
C -->|是| S[ErrStorageInsufficient 配额不足]注意:减少分支的
used_storage >= -delta与增加分支的used_storage <= total_storage - delta都依赖同一条 SQL 内的"读-改-写"原子性(InnoDB 行锁 + 条件更新),这正是 01-user.md 5.8 强调的"让数据库当场读当场改",应用层不再"先读后改"。
5.6.3 配额预检 CheckStorageAvailable(上传前的友好拦截)
CheckStorageAvailable 在上传/分享转存/复制/分片完成之前调用,提供"快速、友好"的拒绝;AddUsedStorage 在落盘写库成功后调用,做"真实记账"。
// 文件: internal/biz/storage.go
func (uc *UserUsecase) CheckStorageAvailable(ctx context.Context, userID uint64, fileSize int64) error {
info, err := uc.GetStorageInfo(ctx, userID) // 读缓存(可能陈旧)
if err != nil {
return err
}
if info.Available < fileSize {
return ErrStorageInsufficient // 413,前端提示"空间不足"
}
return nil
}
调用方示例(上传主流程,internal/biz/file.go):
// 文件: internal/biz/file.go(Upload 关键片段)
func (uc *FileUsecase) Upload(ctx context.Context, userID uint64, name, fileType string, parentID *uint64, data []byte) (*File, error) {
fileSize := int64(len(data))
// 1) 上传前预检配额(基于缓存余量,UX 层快速拒绝)
if uc.userUC != nil {
if err := uc.userUC.CheckStorageAvailable(ctx, userID, fileSize); err != nil {
return nil, err
}
}
// 2) 写对象存储
storagePath, err := uc.storage.Upload(ctx, fileName, bytesToReader(data))
if err != nil {
return nil, err
}
// 3) 写文件元数据
created, err := uc.fileRepo.Create(ctx, file)
if err != nil {
uc.storage.Delete(ctx, storagePath) // 回滚孤儿物理文件
return nil, err
}
// 4) 落库成功后再真实记账
if uc.userUC != nil {
if err := uc.userUC.AddUsedStorage(ctx, userID, fileSize); err != nil {
// 修正点:AddUsedStorage 失败应回滚(删文件 + 删元数据),否则产生"文件在但没扣额度"的漂移
log.Error("file: failed to add used storage after upload", "userID", userID, "err", err)
}
}
return created, nil
}
配额记账的完整生命周期(谁加、谁减):
| 业务动作 | 调用 | 方向 | 说明 |
|---|---|---|---|
| 普通上传 / 秒传 / 分片完成 / 分享转存 / 复制 | AddUsedStorage | +size | 先 CheckStorageAvailable 预检,再落盘记账 |
| 回收站永久删除 | SubUsedStorage | -size | permanentlyDeleteItem 中扣减 |
| 软删除(移入回收站) | 无 | 不变 | 文件仍占空间,配额不释放 |
| 回收站恢复 | 无 | 不变 | 与软删除对称,配额保持,逻辑自洽 |
flowchart LR
Req[上传请求] --> P{CheckStorageAvailable
缓存余量够?}
P -->|否| E[413 空间不足]
P -->|是| W[写对象存储 + 写文件元数据]
W --> A[AddUsedStorage 原子 +size]
A -->|失败| R[回滚 删文件/元数据]
A -->|成功| OK[失效缓存 + 返回]修正点(预检的边界):CheckStorageAvailable 基于 GetStorageInfo 的缓存值,且发生在锁之外。它是 UX 层(不让用户白传完才报错),但不是强制层。强制层是 5.6.2 里 UpdateUsedStorageAtomic 的 DB 上界条件——即便预检因缓存陈旧而误放行,真实记账时也会被 DB 拦下并返回 ErrStorageInsufficient。两者必须同时存在:预检负责体验,DB 条件负责正确。
5.6.4 缓存策略 GetStorageInfo(穿透 / 击穿 / 雪崩 / 前缀 / 失效)
// 文件: internal/biz/user.go(教学版)
const (
cacheKeyUserStorage = "user:storage:%d" // 修正点:缺命名空间前缀
cacheTTLUserStorage = 300 * time.Second // 固定 5 分钟,有雪崩风险
)
func (uc *UserUsecase) GetStorageInfo(ctx context.Context, userID uint64) (*StorageInfo, error) {
if uc.cache != nil {
if val, err := uc.cache.Get(ctx, fmt.Sprintf(cacheKeyUserStorage, userID)); err == nil && val != "" {
var info StorageInfo
if json.Unmarshal([]byte(val), &info) == nil {
return &info, nil // 命中直接返回
}
}
}
// 未命中查库
total, used, err := uc.repo.GetUserStorage(ctx, userID)
if err != nil {
return nil, err
}
available := total - used
if available < 0 {
available = 0
}
var percent float64
if total > 0 {
percent = float64(used) / float64(total) * 100
}
info := &StorageInfo{Total: total, Used: used, Available: available, Percent: percent}
if uc.cache != nil {
if b, err := json.Marshal(info); err == nil {
_ = uc.cache.Set(ctx, fmt.Sprintf(cacheKeyUserStorage, userID), string(b), cacheTTLUserStorage)
// 修正点:set 失败被忽略;且未缓存"空值"
}
}
return info, nil
}
额度变化后如何失效?在 updateStorageWithLock 中:
// 文件: internal/biz/storage.go
func (uc *UserUsecase) updateStorageWithLock(ctx context.Context, userID uint64, delta int64) error {
if delta == 0 {
return nil
}
lockKey := fmt.Sprintf(storageLockKey, userID)
if uc.locker != nil {
if err := uc.locker.Lock(ctx, lockKey); err != nil {
return err
}
defer func() { _ = uc.locker.Unlock(ctx, lockKey) }()
}
if err := uc.repo.UpdateUsedStorageAtomic(ctx, userID, delta); err != nil {
return err
}
// 记账成功后失效缓存,下次读取回源得到新值
if uc.cache != nil {
_ = uc.cache.Delete(ctx, fmt.Sprintf(cacheKeyUserStorage, userID))
}
return nil
}
商用修正点(缓存三大经典问题):
- 缓存穿透(不存在的 userID 反复打 DB):当前对不存在的用户,
GetUserStorage返回ErrUserNotFound,但没有缓存空值,于是每次请求都打到 MySQL。- 修正:缓存一个短 TTL(如 60s)的空标记(如
"null"),GetStorageInfo命中空标记直接返回错误,避免穿透。
- 修正:缓存一个短 TTL(如 60s)的空标记(如
- 缓存击穿(热点用户 key 过期瞬间惊群):大 V 用户的存储信息被高频查询,key 过期的瞬间大量并发同时回源 DB。
- 修正:缓存未命中重建时加 singleflight / 互斥锁,只放一个请求回源,其余等待复用结果。
- 缓存雪崩(大量 key 同时失效):所有 key 都是固定 300s TTL,若批量回源或进程重启后同时写入,会在同一时刻集体过期,DB 瞬时压力陡增。
- 修正:TTL 加随机抖动,例如
300s + rand(-60s, +60s)。
- 修正:TTL 加随机抖动,例如
flowchart TB
G[GetStorageInfo] --> C{缓存命中?}
C -->|命中 含空标记| N[直接返回 防穿透]
C -->|未命中| S[singleflight 抢锁
仅一个回源]
S --> DB[(MySQL)]
DB --> W[写回缓存 TTL+抖动]
W --> R[返回]
U[AddUsedStorage/SubUsedStorage] --> X[记账后 Delete 缓存]
X --> G- 缓存 key 前缀:当前
user:storage:%d没有业务命名空间。在共享 Redis 实例中容易与其他模块/其他服务撞 key。- 修正点:改为
cloud-disk:user:storage:%d(与 01-user.md 的存储缓存前缀保持一致),隔离命名空间。
- 修正点:改为
- 失效 vs 写穿:当前采用"删缓存"策略(delete-after-write)。它简单且与 DB 最终一致;但因删除在 DB 提交之后、且错误被忽略(
_ =),存在一个极短窗口可能读到旧值——对配额正确性无影响(强制在 DB),只影响提示精度。更强做法是写穿(直接SET新值),但需保证并发下写穿值不落后。 - 一致性结论:配额强制在 DB 层(5.6.2),缓存只是加速展示。缓存陈旧最多让用户看到"余量偏多/偏少"几秒,不会导致超卖或扣减丢失——这是"缓存可丢、DB 不可丢"的分层设计。
5.6.5 并发与分布式锁(锁只挡冲突,不挡错误)
biz.Locker 是一个最小分布式锁接口,updateStorageWithLock 在记账前后用 user:storage🔒{id} 串行化同一用户的多次增减:
// 文件: internal/biz/storage.go
const (
storageLockKey = "user:storage🔒%d"
calibrateLockKey = "storage:calibrate:lock"
)
- 作用:把同一用户的并发"上传/删除"串行化,减少
used_storage在应用层被交替修改带来的冲突与缓存抖动。粒度按userID,不同用户之间互不阻塞,粒度合理。 - 边界(关键):
if uc.locker != nil——锁是可选的。教学版若未注入真实 Redis 锁(如locker为nil),则完全没有应用层串行化。即使注入了锁,CheckStorageAvailable也在锁之外执行。 - 结论:分布式锁是"减少冲突的优化",不是正确性的唯一保证。正确性的唯一裁判是 5.6.2 里
UpdateUsedStorageAtomic的 DB 上界条件 + 行锁。商用部署应两者都上:锁降低冲突概率与缓存失效频率,DB 条件兜底绝对正确。这样即便锁服务抖动/未配置,配额也不会被超卖。
5.6.6 监控与可观测(新增,商用必备)
教学版对配额"只记不报":超限只返回错误,没有指标、没有预警、没有给用户的临近提示。商用需补齐:
- 指标暴露(Prometheus):在
UpdateUsedStorageAtomic/GetStorageInfo处上报cloud_disk_user_storage_used_bytes{user_id="123"}cloud_disk_user_storage_percent{user_id="123"}(= used/total*100) 便于 Grafana 看板与告警规则。
- 分级告警:
percent > 80%:预警(提醒用户清理 / 考虑扩容套餐)。percent > 95%:临界告警(near-limit),主动推送通知(站内信 / 邮件 / 短信)提醒"空间即将用尽"。
- 前端提示:在
StorageInfo增加Warning bool/NearLimit bool字段,GetStorageInfo按percent阈值填充,前端据此在界面顶部提示"存储空间已使用 92%,请及时清理"。 - 校准异常告警:
RecalibrateStorage计算出|delta|超过阈值(例如 > 100MB 或 > 总量 5%)时,记一条告警指标——这通常意味着出现了"上传成功但没扣 / 删除成功但没减"的异常路径,需要排查(见 5.6.7)。
flowchart LR
M[用量率指标] --> P{percent?}
P -->|<80%| OK[正常]
P -->|80%~95%| W[预警 提醒清理]
P -->|>95%| C[临界 推送通知]
R[RecalibrateStorage 漂移] -->|超阈值| A[告警 排查异常路径]5.6.7 大用户与分库分表(聚合性能 / 分片后原子性)
RecalibrateStorage 是每日兜底校准:用"实际活跃文件总大小"反算 used_storage,纠正累计偏差。它依赖两个查询:
// 文件: internal/data/user.go
func (r *userRepo) GetActiveFileTotalSize(ctx context.Context, userID uint64) (int64, error) {
var total int64
err := r.db.WithContext(ctx).Model(&model.File{}).
Where("user_id = ? AND status = 0", userID).
Select("COALESCE(SUM(size), 0)").Scan(&total).Error
return total, err
}
// 文件: internal/biz/storage.go
func (uc *UserUsecase) RecalibrateStorage(ctx context.Context, userID uint64) error {
if uc.locker != nil {
uc.locker.Lock(ctx, calibrateLockKey)
defer func() { _ = uc.locker.Unlock(ctx, calibrateLockKey) }()
}
actualUsed, err := uc.repo.GetActiveFileTotalSize(ctx, userID)
if err != nil {
return err
}
_, used, err := uc.repo.GetUserStorage(ctx, userID)
if err != nil {
return err
}
delta := actualUsed - used
if delta != 0 {
uc.repo.UpdateUsedStorageAtomic(ctx, userID, delta) // 用原子更新纠正,避免覆盖并发改动
}
if uc.cache != nil {
_ = uc.cache.Delete(ctx, fmt.Sprintf(cacheKeyUserStorage, userID))
}
return nil
}
商用修正点(大用户 / 海量用户场景):
- 大用户聚合慢:
GetActiveFileTotalSize是SUM(size)全表扫描该用户活跃文件。一个用户若有百万级文件,SUM很慢;而RecalibrateStorage若每日对所有用户跑一遍,整体开销巨大。- 修正:
used_storage本身就是"运行计数器",主路径靠原子更新维护,校准只作兜底。校准应限流、错峰(低峰期)、按用户重要性/漂移风险抽样执行,而不是全量每天跑。
- 修正:
ListAllUserIDs全量载入内存:校准遍历用户时若用ListAllUserIDs一次性Pluck所有 ID,海量用户下会 OOM。- 修正:改为基于
id游标的分页(WHERE id > ? ORDER BY id LIMIT 1000),逐批处理,内存恒定。
- 修正:改为基于
- 分库分表后原子性仍成立:若按
user_id哈希分库分表,单个用户行落在同一分片,UpdateUsedStorageAtomic的WHERE ... AND 行锁在分片内依然原子成立——并发上传的扣减一致性不受分库影响。 - 跨用户汇总需独立通道:运营后台要看"全站已用总存储"等跨用户聚合,不能再靠逐用户
SUM。应维护独立的汇总表 / 物化视图,或离线计算任务,避免在线校准链路被拖垮。 - 校准与在线记账的协作:校准用
UpdateUsedStorageAtomic(delta)而非直接Updates(used=actual),是为了不覆盖校准期间发生的并发上传/删除——delta是增量,原子更新保证与并发改动叠加正确。
自测题与动手练习
自测题(合上书能答出来,才算懂):
biz.Storage接口一共有几个方法?其中哪几个属于"分片上传生命周期"?哪几个方法要求"幂等"?var _ biz.Storage = (*localStorage)(nil)这行编译期断言的作用是什么?如果localStorage漏实现Delete,会发生什么?- 本地存储和 MinIO 在"分片合并"上分别怎么做的?为什么说 MinIO 的
ComposeObject是"服务端合并"、更省带宽? - 什么是"孤儿分片"?它是怎么产生的?
CleanupOrphanChunks用ModTime/LastModified判断过期的逻辑是什么? UpdateUsedStorageAtomic在「减少空间」时为什么要在WHERE里加used_storage >= ??在「增加空间」时教学版漏了什么,会导致什么后果?RowsAffected == 0时要区分哪两种情况?CheckStorageAvailable和UpdateUsedStorageAtomic的 DB 上界条件,哪个是"强制层"、哪个是"UX 层"?为什么光靠预检还不够?GetStorageInfo的缓存 key 当前有什么问题(前缀 / 雪崩)?缓存穿透、击穿、雪崩分别怎么防?额度变化后缓存如何失效?- 分布式锁
user:storage🔒{id}能保证配额不出错吗?为什么我们说"DB 条件才是唯一裁判"? - 软删除、回收站恢复、永久删除分别对
used_storage有什么影响?为什么软删除不释放配额? - 大用户校准时
GetActiveFileTotalSize和ListAllUserIDs各有什么性能隐患?分库分表后UpdateUsedStorageAtomic的原子性还成立吗?
动手练习(建议真做一遍):
- 把
Storage接口的 9 个方法画成一张职责图,标出普通文件操作、分片生命周期、孤儿清理三组,并对照工厂NewStorage想清楚"加一个 OSS 实现要改几处"。 - 给
CleanupOrphanChunks设计一个定时调用(如每天凌晨跑一次),并写一句说明:在你们的业务量下,maxAge设 24 小时还是 7 天更合理,为什么? - 在
UpdateUsedStorageAtomic的增加分支补上WHERE used_storage <= total_storage - ?,写一段并发测试:起 10 个 goroutine 各上传 200MB,用户配额 1GB,验证最终used_storage不会超过 1GB(无超卖)。 - 给
GetStorageInfo的缓存加"空值缓存 + TTL 抖动 + singleflight",并用 ab/hey 模拟"热点用户 key 过期瞬间 100 并发",观察 DB 是否只被命中一次(防击穿)。 - 在
StorageInfo增加NearLimit bool,在GetStorageInfo里当percent > 95时置 true,并在前端(或伪代码)写出"空间即将用尽"的提示分支。
本章小结
- 接口抽象是核心:
biz.Storage用 9 个方法覆盖普通文件操作、分片上传全生命周期、孤儿分片清理,业务层零感知底层;工厂 + Wire 注入让新增存储引擎只需一个case,编译期断言守住"不漏方法"底线。 - 路径分散有讲究:
datePath()/dateKey()让文件天然分散避免单目录爆炸;MinIO 的"目录"只是对象 key 前缀,无真实层级;分片合并本地io.CopyBuffer+Sync,MinIOComposeObject服务端合并更省应用服务器带宽。 - 配额原子更新是底线:
UpdateUsedStorageAtomic用gorm.Expr在同一条 SQL 内"读-改-写",减少分支WHERE used_storage >= ?防扣成负、增加分支必须补WHERE used_storage <= total_storage - ?防并发超卖;RowsAffected == 0要区分"用户不存在"与"配额不足"。 - 预检与 DB 双保险:
CheckStorageAvailable基于缓存做上传前友好拒绝(UX 层),真正强制在 DB 上界条件(正确层);软删除不释放配额、永久删除才SubUsedStorage、恢复对称不变,逻辑自洽。 - 缓存只加速不保正确:
GetStorageInfo「缓存优先 → 查库 → 回写」,记账后失效缓存;商用需补命名空间前缀(cloud-disk:)、TTL 抖动(防雪崩)、空值缓存(防穿透)、singleflight(防击穿)。配额强制在 DB,缓存陈旧不影响正确性。 - 锁只挡冲突:分布式锁串行化同一用户,但可选且不在预检路径上,故不是正确性唯一保证;DB 条件 + 行锁才是最终裁判。
- 可观测与大数据:用量率指标 + 80%/95% 分级告警 + 临近提示是商用必备;大用户校准用
SUM慢、遍历用游标分页,分库后用户行仍在单分片、原子性不变,跨用户汇总走独立通道。