文件批处理服务
日处理约 5 万批次 · 单批几十至四千次 OSS 操作 · 数据不落地
背景与结果
公司内部 OSS 文件批处理服务:业务方提交文件操作指令(复制文件/多文件/多目录、设置 ACL),服务异步执行并回传进度与结果。
建设动机:文件操作时效要求不一(有的需快速完成、有的批处理即可);文件操作 IO 密集且长耗时,在原业务服务上开多线程会挤占原主机资源影响其他业务——独立服务 + 独立机器做物理隔离。
2023.03 - 2023.05(3 个月上线初版),一人独立完成架构设计、开发到部署上线。
结果数据:日处理约 5 万批次;单批次底层 OSS 操作几十到四千次不等(底层 OSS 操作日峰值千万级)。
技术方案
任务模型
- 一次可提交同一批多个指令,支持指定指令先后顺序(同批顺序执行)
- 提交即入库,返回任务序列号;结果通知 + 查询双通道(可选携带通知 URL 回调带序列号,查询接口始终可用)
- 失败策略由调用方配置:同批顺序指令遇失败时中止还是继续——有依赖的批次选中止,无依赖的选继续
- 重试口径:任务失败后立即重试 1 次,再失败则记录失败次数、不再重试——OSS 服务端复制失败多为瞬时抖动(网络/限流),立即重试一次能消化大多数;连续两次失败大概率是持久性问题(权限/文件不存在),重试无意义
- 任务表按年份归档(超期数据定时迁移归档表,热表保当年)
多优先级 Worker 架构(亮点一)
高/中/低三组独立线程组,线程数 / 每次批量获取数 / 轮询间隔 / 重试次数全部可配置——优先级语义贯穿"线程数 × 拉取节奏 × 批量"三个维度,改配置不改代码:
| 优先级 | 线程数 | 每次批量获取 | 轮询间隔 | 定位 |
|---|---|---|---|---|
| 高 | 512 | 4 | 500ms | 需快速完成的任务 |
| 中(默认) | 256 | 3 | 可配 | 常规任务 |
| 低 | 64 | 2 | 2s | 不及时的批处理,慢速消化 |
- 每组独立 Redis 分布式锁:优先级资源级隔离——低优先级积压永不影响高优先级时效
- 弹性扩展:任务处理不完时直接拉机器镜像启动即成集群,多机抢锁天然分摊负载、不重复消费
这个架构解决的是线程池的经典问题:单一线程池 + 优先级队列看似也能分优先级,但低优先级任务积压会占满队列,高优先级任务排队等锁——分组后各组有独立的线程、拉取节奏与锁,隔离是物理的。
旁路子线程 + 队列(亮点二)
进度通道:每执行完一个任务推一条结果进进度队列;专职子线程刷库时按批次分组、成功/失败各自累加,合并成一次 UPDATE——写库次数与批次数挂钩而非任务数挂钩。一个 4000 任务的批次,进度写库从 4000 次降为常数次。进度模型三计数器:总任务数 / 成功任务数 / 失败任务数。
通知通道:结果通知投通知队列,专职子线程发 HTTP——通知绝不在任务处理线程中执行,对端故障不影响任何处理线程(故障域分离)。通知失败梯度重试 3/5/10 分钟(从首次失败起算,间隔可配)。
完整告警闭环
- 通知失败升级告警:重试穷尽 → 该业务方当天通知失败计数累加 → 达阈值调用统一消息服务发短信给相关负责人——管"结果送不出去"(下游故障)
- 积压告警:待处理任务数达阈值 → 短信通知上管理后台查看——管"任务消化不动"(上游洪峰/自身故障),同时正是拉镜像扩容的触发信号:告警 → 看各优先级队列深度 → 决定加机器
复制实现
全部操作 OSS 服务端复制——数据不落地本服务、不经本地中转;服务只发 API 指令,832 个线程全是轻量 IO 等待,单机可扛。