# timeout-center-demo **Repository Path**: bjfengwen/timeout-center-demo ## Basic Information - **Project Name**: timeout-center-demo - **Description**: No description available - **Primary Language**: Java - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-03-18 - **Last Updated**: 2026-03-18 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # timeout-center-demo 这是一个基于 Spring Boot + MySQL(JPA) + Redis ZSet 的轻量级超时中心示例工程。 它实现的不是简单定时任务,而是一套带有以下能力的超时任务调度模型: - 任务持久化 - 分片路由 - 低延迟调度 - 并发消费 - 失败重试 - prepare 悬挂恢复 - 死信兜底 --- ## 1. 项目解决什么问题 很多业务都需要“未来某个时间点触发一次处理”,例如: - 订单超时取消 - 支付超时关闭 - 预约超时释放 - 营销资格到期处理 如果只用 DB 轮询,实时性通常不够;如果只用 Redis,又缺少可靠状态与恢复能力。 这个 demo 的核心思路是: - **DB 保存任务权威状态** - **Redis 负责低延迟调度** 因此它兼顾了: - 可恢复 - 可追踪 - 执行时效高 --- ## 2. 核心业务流程 ### 2.1 创建任务 入口: - `src/main/java/com/example/timeoutcenter/api/TimeoutJobController.java` - `src/main/java/com/example/timeoutcenter/domain/service/TimeoutJobService.java` 创建流程: 1. 调用 `POST /jobs` 2. `TimeoutJobService.createJob(...)` 生成 `jobId` 3. 根据 `topic + slotBasis + jobId` 计算目标 `slot` 4. 先把任务写入 DB 5. 再把任务写入 Redis `store` 队列,score = `actionTime` ### 为什么同时写 DB 和 Redis? #### DB 的职责 - 保存任务元数据 - 保存任务最终状态 - 保存重试次数 - 提供恢复、排查、审计基础 #### Redis 的职责 - 按时间排序 - 让 worker 快速拿到到期任务 - 适合高频轮询和低延迟调度 一句话: > DB 是权威状态层,Redis 是调度层。 --- ### 2.2 slot 路由 相关类: - `src/main/java/com/example/timeoutcenter/domain/service/SlotRouter.java` - `src/main/java/com/example/timeoutcenter/domain/service/TopicRegistry.java` - `src/main/java/com/example/timeoutcenter/config/TimeoutCenterProperties.java` 系统支持多个 topic,每个 topic 可以配置多个 slot。 创建任务时,不是直接丢到一个总队列,而是: - 先确定 topic - 再根据 `slotBasis` 做 hash - 路由到该 topic 下某个固定 slot 路由算法: 1. 优先使用 `slotBasis` 2. 如果没传,就退化为 `jobId` 3. 用 CRC32 计算 hash 4. 用 `hash & (slotAmount - 1)` 得到 slot ### slot 的意义 slot 的核心作用: - 分片 - 降低竞争 - 支撑并行消费 - 让同一类业务 key 稳定落到同一分片 --- ### 2.3 Redis 三类队列 核心类: - `src/main/java/com/example/timeoutcenter/infra/redis/RedisTimeoutQueueRepository.java` - `src/main/java/com/example/timeoutcenter/infra/redis/LuaScripts.java` 系统核心有三类 Redis ZSet 队列: #### 1)store 队列 待执行主队列。 特点: - 新任务先进入这里 - score = `actionTime` - worker 只抢占 `score <= now` 的任务 #### 2)prepare 队列 已被 worker 抢占但尚未确认完成的队列。 特点: - worker 抢到任务后,会先把任务从 `store` 移到 `prepare` - score = `prepareExpireAt` - 如果超过 prepare 超时时间还没 ack,说明可能处理悬挂了 #### 3)dead 队列 死信队列。 特点: - 超过最大重试次数后进入 - 不再自动处理 - 用于人工排查、补偿 --- ### 2.4 worker 消费主链路 核心类: - `src/main/java/com/example/timeoutcenter/scheduler/TimeoutWorker.java` worker 启动后会做两件事: 1. 按 `workerCount` 启动多个消费线程 2. 每个线程只轮询自己负责的 slot slot 分工逻辑: - `slot % workerCount == workerIndex` 也就是说: - 每个 slot 都会有一个固定 worker 负责 - 避免所有 worker 竞争所有 slot ### 消费流程 1. worker 扫描自己负责的 slot 2. 调用 `claimDueJob(...)` 3. 用 Lua 把一个到期 job 从 `store` 原子迁移到 `prepare` 4. 加载 DB 中的任务状态 5. 如果状态仍为 `PENDING`,则执行对应 `TimeoutJobProcessor` 6. 根据执行结果进入不同分支: - 成功:ack,任务完成 - 失败且未超限:rollback,回到 store 等待重试 - 失败且超限:转 dead --- ### 2.5 为什么 claim 要用 Lua 脚本 核心脚本: - `src/main/java/com/example/timeoutcenter/infra/redis/LuaScripts.java` Lua 脚本做的事: 1. 从 `store` 中找一个已到期任务 2. 从 `store` 删除 3. 加入 `prepare` 4. 返回 jobId 必须用 Lua 的原因: 如果拆成多步 Java 操作: - worker A 查到了一个 job - worker B 也查到了同一个 job - 两边都觉得自己抢到了 - 就会重复消费 Lua 的意义是把: - 查找 - 删除 - 移入 prepare 放在 Redis 里原子执行,避免并发重复抢占。 --- ### 2.6 成功、失败、重试、死信 处理逻辑在: - `src/main/java/com/example/timeoutcenter/scheduler/TimeoutWorker.java` #### 成功路径 - processor 执行业务逻辑成功 - DB 状态改为 `PROCESSED` - 从 `prepare` 删除 链路: `store -> prepare -> processed + ack` #### 失败重试路径 - processor 抛异常 - `retryCount + 1` - 如果未超过 `maxRetries` - 计算下一次重试时间 - 从 `prepare` 回滚到 `store` 链路: `store -> prepare -> rollback -> store` #### 死信路径 - 失败次数超过上限 - DB 状态改为 `DEAD` - 从 `prepare` 转入 `dead` 链路: `store -> prepare -> dead` --- ### 2.7 prepare 恢复流程 核心类: - `src/main/java/com/example/timeoutcenter/scheduler/PrepareRecoveryWorker.java` 恢复线程的职责是: > 扫描 prepare 队列里“抢到了但迟迟没 ack”的任务。 这类情况通常代表: - worker 处理过程中宕机 - 线程异常退出 - 业务执行卡死 - 已经抢占但没有机会完成 ack 恢复流程: 1. 周期性扫描所有 topic/slot 的 prepare 队列 2. 找到 `score <= now` 的 prepare 任务 3. 对每个过期任务执行: - `retryCount + 1` - 如果超过最大重试次数:转 dead - 否则:回滚到 store,等待后续再次消费 因此 prepare 队列实际上承担了一个“处理中间态 + 失败恢复锚点”的作用。 --- ### 2.8 取消流程 逻辑在: - `src/main/java/com/example/timeoutcenter/domain/service/TimeoutJobService.java` 取消任务时会: 1. DB 状态改成 `CANCELED` 2. 从 `store` 删除 3. 从 `prepare` 删除 这样做是为了避免: - 任务还在待执行队列里继续被消费 - 或者已经被抢占后又被恢复线程重新投递 --- ## 3. 业务时序图 ### 3.0 Mermaid 图 #### 核心业务时序图 ```mermaid sequenceDiagram participant C as Client participant API as TimeoutJobController participant SVC as TimeoutJobService participant DB as MySQL/DB participant RStore as Redis Store ZSet participant W as TimeoutWorker participant P as TimeoutJobProcessor participant RPrepare as Redis Prepare ZSet participant REC as PrepareRecoveryWorker participant RDead as Redis Dead ZSet C->>API: POST /jobs API->>SVC: createJob(...) SVC->>SVC: 生成 jobId / 计算 slot SVC->>DB: 保存任务(status=PENDING) SVC->>RStore: schedule(jobId, actionTime) API-->>C: 返回任务信息 loop 到期轮询 W->>RStore: claimDueJob(now) RStore->>RPrepare: Lua原子迁移 store -> prepare RPrepare-->>W: 返回 jobId W->>DB: requireJob(jobId) alt 状态不是 PENDING W->>RPrepare: ack(remove prepare) else 状态为 PENDING W->>P: process(job) alt 处理成功 W->>DB: markProcessed(jobId) W->>RPrepare: ack(remove prepare) else 处理失败且未超限 W->>DB: incrementRetryCount(jobId) W->>RStore: rollback(prepare -> store, retryTime) else 处理失败且超限 W->>DB: markDead(jobId) W->>RDead: moveToDead(prepare -> dead) end end end loop prepare恢复扫描 REC->>RPrepare: listExpiredPrepareJobs(now) alt prepare任务超时且未超限 REC->>DB: incrementRetryCount(jobId) REC->>RStore: rollback(prepare -> store, retryTime) else prepare任务超时且超限 REC->>DB: markDead(jobId) REC->>RDead: moveToDead(prepare -> dead) end end ``` #### 系统架构图 ```mermaid flowchart TB Client[调用方 / Client] --> API[TimeoutJobController\nAPI层] API --> SVC[TimeoutJobService\n领域服务层] SVC --> Router[SlotRouter\n分片路由] SVC --> Registry[TopicRegistry\nTopic元信息] SVC --> DB[(MySQL / timeout_job)] SVC --> RedisRepo[RedisTimeoutQueueRepository\nRedis调度层] RedisRepo --> Store[(store ZSet)] RedisRepo --> Prepare[(prepare ZSet)] RedisRepo --> Dead[(dead ZSet)] RedisRepo --> Lua[LuaScripts\n原子抢占] Worker[TimeoutWorker\n消费线程池] --> RedisRepo Worker --> SVC Worker --> Processor[TimeoutJobProcessor\n业务处理扩展点] Recovery[PrepareRecoveryWorker\n恢复线程] --> RedisRepo Recovery --> SVC Config[TimeoutCenterProperties\n配置中心] --> Worker Config --> Recovery Config --> Registry ``` ### 3.1 正常成功链路 ```text Client | | POST /jobs v TimeoutJobController | v TimeoutJobService |-- 生成 jobId |-- 路由 slot |-- 写 DB(timeout_job, status=PENDING) |-- 写 Redis store(score=actionTime) v 返回任务信息 ... 到达 actionTime 后 ... TimeoutWorker |-- 扫描负责的 slot |-- Lua claim: store -> prepare |-- 读取 DB 状态 |-- 调用 TimeoutJobProcessor |-- DB status -> PROCESSED |-- ack: remove prepare v 任务完成 ``` ### 3.2 失败重试链路 ```text TimeoutWorker |-- claim: store -> prepare |-- process(job) 失败 |-- retryCount + 1 |-- 未超过 maxRetries |-- rollback: prepare -> store(score=retryTime) v 等待退避后再次被 worker 消费 ``` ### 3.3 prepare 恢复链路 ```text TimeoutWorker |-- claim: store -> prepare |-- 处理中异常退出 / 卡死 / 未 ack v prepare 中残留任务 PrepareRecoveryWorker |-- 扫描 prepare 过期任务 |-- retryCount + 1 |-- 未超限: rollback -> store |-- 超限: move -> dead ``` ### 3.4 死信链路 ```text 任务多次失败 | | retryCount > maxRetries v DB status = DEAD Redis: prepare -> dead v 等待人工排查或补偿处理 ``` --- ## 4. 系统架构梳理 ### 4.1 总体结构 ```text 调用方 / Client | v TimeoutJobController API 层 | v TimeoutJobService 领域服务层 / \ v v DB持久化层 Redis调度层 (JPA) (ZSet + Lua) | | | | +-------> TimeoutWorker <-------+ | v TimeoutJobProcessor 另外一条后台恢复链路: PrepareRecoveryWorker ``` --- ### 4.2 分层职责 #### 1)API 层 类: - `TimeoutJobController` 职责: - 对外提供创建、取消、查询接口 - 不承载核心调度逻辑 #### 2)领域服务层 类: - `TimeoutJobService` - `SlotRouter` - `TopicRegistry` 职责: - 创建/取消任务 - 管理 DB 权威状态 - 决定 slot 路由 - 管理 topic 元信息 #### 3)DB 持久化层 类: - `TimeoutJobEntity` - `TimeoutJobRepository` 职责: - 保存任务元数据 - 保存状态、slot、retryCount、actionTime - 支撑恢复、查询、审计 #### 4)Redis 调度层 类: - `RedisTimeoutQueueRepository` - `LuaScripts` 职责: - 维护 store / prepare / dead 三类 ZSet - 维护 claim / ack / rollback / dead 迁移 - 用 Lua 保证 claim 原子性 #### 5)调度执行层 类: - `TimeoutWorker` - `PrepareRecoveryWorker` 职责: - worker:消费到期任务 - recovery:恢复 prepare 超时任务 #### 6)业务处理扩展层 类: - `TimeoutJobProcessor` 职责: - 按 topic 选择处理器 - 承载真正的业务执行逻辑 - 与调度框架解耦 --- ## 5. 配置驱动点 配置类: - `src/main/java/com/example/timeoutcenter/config/TimeoutCenterProperties.java` 关键配置: - `workerCount`:消费线程数 - `batchSize`:每次扫描数量 - `pollInterval`:worker 空轮询间隔 - `recoveryInterval`:恢复线程扫描间隔 - `prepareTimeout`:prepare 超时阈值 - `retryBackoff`:重试退避时长 - `maxRetries`:最大重试次数 - `workersEnabled`:是否开启后台 worker - `topics`:topic 及 slotAmount 配置 说明: 这个项目是一个比较典型的**配置驱动型调度框架**。 --- ## 6. 集成测试覆盖了哪些核心场景 测试文件: - `src/test/java/com/example/timeoutcenter/TimeoutCenterIntegrationTest.java` 主要验证了 4 条主链路: 1. **创建并成功消费** - 任务正常创建 - 到时成功执行 - 状态变为 `PROCESSED` 2. **失败重试后成功** - 前几次失败 - rollback 回 store - 后续重试成功 3. **prepare 过期恢复** - 人工构造过期 prepare 任务 - recovery 线程把它恢复回 store 4. **超限进入死信** - 多次失败后转入 `dead` - DB 状态变 `DEAD` 这 4 个测试基本把系统主生命周期覆盖完整了。 --- ## 7. 面试讲解版总结 如果面试时介绍这个项目,可以按下面方式讲: ### 一句话介绍 这是一个基于 **Spring Boot + JPA + Redis ZSet** 实现的轻量级超时任务中心,支持低延迟调度、失败重试、prepare 恢复和死信兜底。 ### 架构亮点 1. **DB + Redis 双层设计** - DB 存权威状态 - Redis 做高性能时间调度 2. **slot 分片模型** - topic 下配置多个 slot - 任务通过 hash 路由到固定 slot - worker 按 slot 静态分工,降低竞争 3. **三段式生命周期队列** - `store`:待执行 - `prepare`:已抢占未确认 - `dead`:终态失败 4. **Lua 原子 claim** - 保证从 `store` 到 `prepare` 的迁移原子化 - 避免并发 worker 重复消费 5. **恢复机制完善** - 如果 worker 抢到任务后挂了,任务不会丢 - prepare 超时后由恢复线程重新回投或转死信 ### 适用场景 - 订单超时关闭 - 支付超时取消 - 延时回调 - 资源预约超时释放 - 任意“未来时间点触发”的业务任务 --- ## 8. 核心设计思想总结 这个 demo 最值得记住的是 6 点: 1. **DB 保状态,Redis 保调度** 2. **slot 分片支撑并发消费** 3. **store / prepare / dead 表达完整生命周期** 4. **Lua claim 保证并发安全** 5. **失败可重试,带退避时间** 6. **prepare 超时可恢复,避免任务悬挂丢失** 如果用一句话概括整个系统: > 这是一个支持分片路由、原子抢占、失败重试、悬挂恢复和死信兜底的轻量级超时调度中心。