文件批处理服务

公司内部 OSS 文件批处理服务——多优先级 Worker 架构与旁路队列的线程工程实践

日处理约 5 万批次 · 单批几十至四千次 OSS 操作 · 数据不落地

JavaSpring BootMySQLRedis阿里云 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 等待,单机可扛。