- Introduced a new section on cross-domain collaboration and aggregation, detailing decision-making processes, contract module usage for cross-domain reads, and domain events for writes. - Added guidelines for parallel aggregation using a dedicated thread pool and context propagation. - Established rules for transaction boundaries, idempotency, optimistic locking, scheduled tasks, and caching strategies in a concurrent environment. - Included examples and best practices for implementing these concepts in the application.
307 lines
17 KiB
Markdown
307 lines
17 KiB
Markdown
# 12. 并发、事务与定时任务
|
||
|
||
## 决策
|
||
|
||
事务边界统一放在 application 层;写接口用幂等键防重复提交;并发冲突用 `@Version` 乐观锁 + 有限重试;定时任务在多副本下用 ShedLock 保证只跑一次;缓存暂不引入 Redis,只用 Caffeine 本地缓存,且严格限定在"能容忍副本间不一致"的数据上。
|
||
|
||
## 一、事务边界
|
||
|
||
### 规则:事务开在 application 层的用例方法上
|
||
|
||
```kotlin
|
||
@Service
|
||
class StoreSwitchAppService(...) {
|
||
|
||
@Transactional // ← 事务边界在这里
|
||
fun switchStore(userId: Long, targetStoreId: Long): StoreContext { ... }
|
||
}
|
||
```
|
||
|
||
- **不在 Controller 上开事务**:Controller 属于 api 层,开事务等于把 HTTP 序列化过程也圈进事务里,事务被无谓拉长。
|
||
- **不在 Repository 方法上开事务**:一个用例往往包含多次写,各自开事务就没有原子性可言了。
|
||
- **只读查询加 `@Transactional(readOnly = true)`**:Hibernate 会跳过脏检查(dirty checking),减少一次全量快照比对;对只读为主的查询接口是免费的性能收益。
|
||
|
||
### 事务里绝对不能做的三件事
|
||
|
||
1. **调外部 HTTP 接口**。F6 慢一点,数据库连接和行锁就被一起占着不放,几个请求就能把连接池打满(见 [03-persistence.md](./03-persistence.md) 的池子容量算法)。外部调用要么放在事务开始前,要么放在提交之后(`AFTER_COMMIT` 事件,见 [11-cross-domain-collaboration.md](./11-cross-domain-collaboration.md))。
|
||
2. **`Thread.sleep` / 等待用户输入 / 等锁**。同上。
|
||
3. **catch 掉异常却继续用同一个事务**。Spring 默认在 `RuntimeException` 时把事务标记为 rollback-only,这之后再做任何写操作,最终提交时都会抛 `UnexpectedRollbackException`——而且报错点离真正的错误现场很远,非常难查。
|
||
|
||
### 自调用失效:Spring 事务最经典的坑
|
||
|
||
```kotlin
|
||
@Service
|
||
class OrderAppService {
|
||
fun createBatch(items: List<Item>) {
|
||
items.forEach { create(it) } // ← 这里的 @Transactional 完全不生效
|
||
}
|
||
|
||
@Transactional
|
||
fun create(item: Item) { ... }
|
||
}
|
||
```
|
||
|
||
`@Transactional` 靠 AOP 代理实现,**同一个类内部的方法调用不经过代理**,注解形同虚设。这一条对 `@Async`、`@Cacheable`、Resilience4j 的注解全部适用。解法是把被调方法挪到另一个 bean 里,让调用真的穿过代理。
|
||
|
||
### 传播行为:只用这两个
|
||
|
||
| 传播行为 | 什么时候用 |
|
||
| --- | --- |
|
||
| `REQUIRED`(默认) | 绝大多数情况。有事务就加入,没有就新建 |
|
||
| `REQUIRES_NEW` | 必须独立提交/回滚的场景:`AFTER_COMMIT` 事件监听器、审计记录、失败次数累加(登录失败计数必须在认证失败回滚后仍然保留) |
|
||
|
||
其余传播行为(`NESTED`、`SUPPORTS`、`MANDATORY`…)在这套系统里没有需要它们的场景,用了只会增加理解成本。真遇到时先怀疑是不是边界划错了。
|
||
|
||
### 隔离级别
|
||
|
||
统一 `READ COMMITTED`(在 [03-persistence.md](./03-persistence.md) 里通过 `transaction-isolation` 全局配置,覆盖 MySQL 默认的 `REPEATABLE READ`),不在方法上单独指定。需要更强一致性的地方用显式锁(`SELECT ... FOR UPDATE`,JPA 的 `@Lock(PESSIMISTIC_WRITE)`)或乐观锁,而不是靠调隔离级别——后者影响面是整个连接,副作用难以预料。
|
||
|
||
## 二、幂等
|
||
|
||
### 哪些接口需要
|
||
|
||
客户端在弱网下会重试(见客户端 `05-networking.md`),用户也会连点两次。**所有会产生副作用且重复执行会出问题的写接口**都需要幂等保护:换票、切店、任何创建类操作。
|
||
|
||
天然幂等的不需要额外处理:`GET`、把状态设为某个确定值的更新(`status = INACTIVE`)、按主键的删除。
|
||
|
||
### 方案:客户端生成幂等键
|
||
|
||
```
|
||
POST /api/v1/webview/tickets
|
||
Idempotency-Key: 7f3a9c1e-... # 客户端生成的 UUID,重试时复用同一个值
|
||
```
|
||
|
||
```kotlin
|
||
// platform/platform-web/src/main/kotlin/.../platform/web/idempotency/IdempotencyGuard.kt
|
||
// 放 platform 而不是某个 domain:换票、切店、创建类操作分散在多个 domain,
|
||
// 而 domain 之间不能互相依赖——公共能力只能住在 platform-*(见 01-project-structure.md)。
|
||
@Service
|
||
class IdempotencyGuard(private val recordRepository: IdempotencyRecordRepository) {
|
||
|
||
/**
|
||
* 同一个 key 在有效期内只会真正执行一次;重复请求直接返回首次的结果。
|
||
*/
|
||
@Transactional
|
||
fun <T> execute(key: String, userId: Long, block: () -> T): T {
|
||
val existing = recordRepository.findByKeyAndUserId(key, userId)
|
||
if (existing != null) return deserialize(existing.response)
|
||
|
||
val result = block()
|
||
// 唯一索引兜底:两个并发请求同时走到这里,第二个会因为 uk_idem_key_user 冲突而失败
|
||
recordRepository.save(IdempotencyRecordEntity(key, userId, serialize(result)))
|
||
return result
|
||
}
|
||
}
|
||
```
|
||
|
||
```sql
|
||
-- 建在 platform 共用的库里(和下面的 shedlock 表同库),不属于任何 domain
|
||
create table idempotency_record (
|
||
id bigint not null auto_increment,
|
||
idem_key varchar(64) not null,
|
||
user_id bigint not null,
|
||
response text not null,
|
||
created_at datetime(6) not null,
|
||
primary key (id),
|
||
unique key uk_idem_key_user (idem_key, user_id)
|
||
) engine = InnoDB default charset = utf8mb4 collate = utf8mb4_0900_ai_ci;
|
||
```
|
||
|
||
三个要点:
|
||
|
||
- **唯一索引是真正的保证,代码里的"先查再写"不是**。两个并发请求会同时查到"没有记录"然后同时执行——只有数据库的唯一约束能挡住。捕获 `DataIntegrityViolationException` 后重新读一次已有结果即可。
|
||
- **幂等键要带 `userId`**:只用 key 做唯一索引意味着不同用户的 key 会互相冲撞,而 key 是客户端生成的,我们不控制它的全局唯一性。
|
||
- **记录要定期清理**(见下面定时任务一节),保留期取"客户端最长可能重试的窗口",比如 24 小时,不需要永久保存。
|
||
|
||
## 三、并发冲突:乐观锁
|
||
|
||
### 用 `@Version`,不用悲观锁
|
||
|
||
```kotlin
|
||
// domains/identity-store/src/main/kotlin/.../identitystore/infrastructure/persistence/StoreEntity.kt
|
||
@Entity
|
||
class StoreEntity : VersionedEntity() { // VersionedEntity 带 @Version,见 03-persistence.md
|
||
var name: String = ""
|
||
}
|
||
```
|
||
|
||
更新时如果 version 已经被别人改过,Hibernate 抛 `ObjectOptimisticLockingFailureException`,[06-api-design.md](./06-api-design.md) 的 `GlobalExceptionHandler` 会把它转成 `10409 CONFLICT`。
|
||
|
||
选乐观锁而不是 `SELECT ... FOR UPDATE`:这套系统是典型的低冲突场景(同一门店的同一条记录被两个人同时改的概率很低),悲观锁的代价是每次读都要持锁,把并发度砍掉换一个几乎用不上的保证。
|
||
|
||
### 什么时候自动重试,什么时候返回 409
|
||
|
||
| 场景 | 处理 |
|
||
| --- | --- |
|
||
| 用户提交的业务更新(改门店信息) | **返回 409**,让用户看到"数据已被他人修改,请刷新后重试"。自动重试会静默覆盖别人的修改 |
|
||
| 内部的计数器/状态推进(登录失败次数、票据状态流转) | **自动重试**,用户不需要知道内部发生了冲突 |
|
||
|
||
```kotlin
|
||
@Retryable(
|
||
retryFor = [ObjectOptimisticLockingFailureException::class],
|
||
maxAttempts = 3,
|
||
backoff = Backoff(delay = 50, multiplier = 2.0, random = true), // 抖动,避免两方同步重试同步再撞
|
||
)
|
||
@Transactional
|
||
fun incrementFailedAttempts(userId: Long) { ... }
|
||
```
|
||
|
||
**重试必须在事务外层**——`@Retryable` 要包住 `@Transactional`,因为冲突发生时事务已经标记回滚,必须开一个全新的事务重新读、重新算、重新写。在事务内部重试是无效的(而且会撞上 rollback-only)。由于两个注解在同一个方法上时代理顺序容易搞错,稳妥做法是把重试和事务拆到两个 bean 上:外层 bean 负责 `@Retryable`,内层 bean 负责 `@Transactional`。
|
||
|
||
## 四、定时任务在多副本下的重复执行
|
||
|
||
### 问题
|
||
|
||
[04-security-auth.md](./04-security-auth.md) 里有一个清理过期 refresh token 的 `@Scheduled` 任务,加上上面幂等记录的清理任务。**`@Scheduled` 在每个 Pod 上都会独立执行**——2 个副本就是每次跑 2 遍。清理任务跑两遍问题不大(删除是幂等的),但只要出现一个"发通知""生成对账单"式的任务,重复执行就是事故。
|
||
|
||
这一条在原来的文档里完全没有提到,而 [09-build-deploy.md](./09-build-deploy.md) 明确配了 `replicas: 2`——也就是说按现有文档实施,上线当天就是重复执行状态。
|
||
|
||
### 方案:ShedLock
|
||
|
||
```groovy
|
||
implementation 'net.javacrumbs.shedlock:shedlock-spring:6.9.2'
|
||
implementation 'net.javacrumbs.shedlock:shedlock-provider-jdbc-template:6.9.2'
|
||
```
|
||
|
||
```sql
|
||
-- 放在 platform 共用的库里,不属于任何 domain
|
||
create table shedlock (
|
||
name varchar(64) not null,
|
||
lock_until datetime(6) not null,
|
||
locked_at datetime(6) not null,
|
||
locked_by varchar(255) not null,
|
||
primary key (name)
|
||
) engine = InnoDB default charset = utf8mb4 collate = utf8mb4_0900_ai_ci;
|
||
```
|
||
|
||
```kotlin
|
||
// platform/platform-persistence/src/main/kotlin/.../platform/persistence/scheduling/SchedulingConfig.kt
|
||
@Configuration
|
||
@EnableScheduling
|
||
@EnableSchedulerLock(defaultLockAtMostFor = "PT10M")
|
||
class SchedulingConfig {
|
||
@Bean
|
||
fun lockProvider(dataSource: DataSource): LockProvider = JdbcTemplateLockProvider(
|
||
JdbcTemplateLockProvider.Configuration.builder()
|
||
.withJdbcTemplate(JdbcTemplate(dataSource))
|
||
.usingDbTime() // 用数据库时间而不是各 Pod 的本地时间,避免时钟漂移导致锁失效
|
||
.build(),
|
||
)
|
||
}
|
||
|
||
// domains/identity-store/src/main/kotlin/.../identitystore/application/RefreshTokenCleanupJob.kt
|
||
// 定时任务放 application 层:它就是一个由时钟而不是 HTTP 请求触发的用例。
|
||
@Component
|
||
class RefreshTokenCleanupJob(private val repository: RefreshTokenRepository, private val clock: Clock) {
|
||
|
||
@Scheduled(cron = "0 17 3 * * *") // 每天凌晨 3:17,见下方"为什么不用整点"
|
||
@SchedulerLock(name = "refreshTokenCleanup", lockAtLeastFor = "PT1M", lockAtMostFor = "PT10M")
|
||
fun cleanup() {
|
||
val deleted = repository.deleteExpiredBefore(clock.instant())
|
||
LoggerFactory.getLogger(javaClass).info("清理过期 refresh token,删除 {} 条", deleted)
|
||
}
|
||
}
|
||
```
|
||
|
||
两个参数的含义容易混:
|
||
|
||
- **`lockAtMostFor`**:锁的最长持有时间,防死锁。持锁的 Pod 被 kill 掉时,锁不会自动释放(数据库里的记录还在),到这个时间才失效。**必须显著大于任务的正常执行时间**,否则任务还没跑完锁就过期了,另一个 Pod 会同时开跑。
|
||
- **`lockAtLeastFor`**:锁的最短持有时间。防的是"任务执行极快 + 各 Pod 时钟有偏差"导致同一个调度点被跑两次。
|
||
|
||
**`usingDbTime()` 不能省**:不加的话 ShedLock 用各个 Pod 的本地时间写锁,Pod 之间有几秒时钟漂移就可能同时抢到锁——那这套机制就白配了。
|
||
|
||
### 为什么 cron 不用整点
|
||
|
||
`0 0 3 * * *` 这种整点时间,全世界所有系统的定时任务都挤在同一秒——数据库、外部依赖、监控在那一刻集体尖峰。错开几分钟(`0 17 3 * * *`)没有任何业务代价,但能把这个尖峰摊平。
|
||
|
||
### 定时任务的其他约定
|
||
|
||
- **必须记日志**:开始、结束、处理条数。没有日志的定时任务出问题时你连"它有没有跑"都不知道。
|
||
- **必须限制单次处理量**:`delete from ... where expires_at < ? limit 1000` 分批删,不要一条 SQL 删几百万行——那会长时间持有行锁并撑爆 binlog。
|
||
- **必须自己兜住异常**:`@Scheduled` 方法抛异常只会被 Spring 记一条日志,任务本身不会重试,但下次调度照常。如果失败需要告警,自己 catch 后打 `ERROR`(见 [08-observability.md](./08-observability.md) 的告警规则)。
|
||
- **任务执行情况应该有指标**:至少一个"上次成功执行时间",用于告警"某个任务已经 3 天没成功跑过了"——这类静默失败光看错误日志是发现不了的。
|
||
|
||
### 替代方案:K8s CronJob
|
||
|
||
对于"跑一次就结束、不需要常驻"的任务(比如大表数据归档),更合适的做法是 K8s `CronJob` 起一个独立 Pod:天然只跑一份,不需要 ShedLock,也不会和在线请求抢应用的线程和连接池。代价是要额外维护一套镜像入口和清单。
|
||
|
||
判断标准:**任务需要用到应用内的业务逻辑 → `@Scheduled` + ShedLock;任务本质是一段独立的数据操作 → CronJob**。
|
||
|
||
## 五、缓存:只用 Caffeine 本地缓存
|
||
|
||
### 为什么暂不引入 Redis
|
||
|
||
引入 Redis 意味着多一个需要部署、监控、备份、排障的有状态组件,以及一整套新的失败模式(连接抖动、大 key、缓存穿透/雪崩)。当前的数据量和并发量还远没有到需要它的程度,而 [04-security-auth.md](./04-security-auth.md) 的 refresh token 已经明确落库、[11-cross-domain-collaboration.md](./11-cross-domain-collaboration.md) 的事件也不需要外部存储——没有哪个场景是非它不可的。
|
||
|
||
### Caffeine 的适用边界(这是重点)
|
||
|
||
本地缓存的本质特征是:**每个 Pod 各缓存一份,副本之间必然不一致,且无法主动失效其他副本的缓存**。所以只能缓存满足这两个条件的数据:
|
||
|
||
1. 变更频率极低;
|
||
2. **短时间内读到旧值不会造成业务错误**。
|
||
|
||
| 数据 | 能不能用本地缓存 | 说明 |
|
||
| --- | --- | --- |
|
||
| 菜单/权限元数据定义 | 可以 | 几乎不变,改了之后延迟几分钟生效可接受 |
|
||
| 门店基础信息(名称、编码) | 可以 | 同上 |
|
||
| 字典/枚举配置 | 可以 | 同上 |
|
||
| **用户对门店的可访问权限** | **不可以** | 收回权限后必须立即生效,这是安全边界(见 [04-security-auth.md](./04-security-auth.md)) |
|
||
| **WebView 票据状态** | **不可以** | 作废必须立即生效 |
|
||
| **任何用户维度的业务数据** | **不可以** | 用户在 Pod A 改的数据,请求打到 Pod B 会读到旧值 |
|
||
|
||
```kotlin
|
||
// platform/platform-persistence/src/main/kotlin/.../platform/persistence/cache/CacheConfig.kt
|
||
@Configuration
|
||
@EnableCaching
|
||
class CacheConfig {
|
||
@Bean
|
||
fun cacheManager(): CacheManager = CaffeineCacheManager().apply {
|
||
setCaffeine(
|
||
Caffeine.newBuilder()
|
||
.maximumSize(1_000) // 必须设上限,否则就是内存泄漏
|
||
.expireAfterWrite(Duration.ofMinutes(10)), // 必须设过期,这是副本间最终一致的唯一保证
|
||
)
|
||
}
|
||
}
|
||
|
||
// domains/identity-store/.../identitystore/application/StoreQueryServiceImpl.kt
|
||
// 缓存加在 application 层的查询方法上,返回的是契约模型而不是 Entity——
|
||
// 缓存里躺着一个游离态的 JPA Entity 是另一类难查的问题(见 02-layering.md)。
|
||
@Cacheable(cacheNames = ["storeBasicInfo"], key = "#storeId")
|
||
fun findStoreBasicInfo(storeId: Long): StoreInfo? { ... }
|
||
```
|
||
|
||
`maximumSize` 和 `expireAfterWrite` 都是**强制**的,不是可选优化:没有 `maximumSize` 的缓存就是一个慢速内存泄漏;没有 `expireAfterWrite` 的本地缓存永远不会和其他副本收敛。
|
||
|
||
**用 `@Cacheable` 时同样注意自调用失效**——它和 `@Transactional` 一样走代理,同类内部调用不生效。
|
||
|
||
### 什么时候升级到 Redis
|
||
|
||
- 需要**跨副本立即失效**某个缓存;
|
||
- 需要跨副本共享状态(分布式限流计数、在线用户数);
|
||
- 缓存数据量大到单 Pod 内存放不下。
|
||
|
||
到那时 Redis 的取舍是清楚的,现在提前引入只是提前承担成本。
|
||
|
||
## 关键规则
|
||
|
||
- 事务边界在 application 层的用例方法上;事务内不做 HTTP 调用、不 sleep。
|
||
- 注意自调用失效:`@Transactional` / `@Async` / `@Cacheable` / Resilience4j 注解在同类内部调用时全部不生效。
|
||
- 有副作用的写接口用 `Idempotency-Key` + 唯一索引做幂等;唯一索引才是保证,"先查再写"不是。
|
||
- 并发冲突用 `@Version` 乐观锁:用户提交的更新返回 `10409`,内部状态推进自动重试(重试包在事务外层)。
|
||
- 所有 `@Scheduled` 任务必须加 `@SchedulerLock`,`LockProvider` 必须 `usingDbTime()`;cron 时间避开整点。
|
||
- 只用 Caffeine 本地缓存,且只缓存"读到旧值不会出错"的数据;`maximumSize` 和 `expireAfterWrite` 强制配置。
|
||
|
||
## 待补充
|
||
|
||
- 幂等记录的具体保留期(需要和客户端确认最长重试窗口)。
|
||
- 是否需要接口级限流(当前只有对下游的舱壁,没有对上游的限流)。
|
||
|
||
## 参考链接
|
||
|
||
- [Spring: Transaction Management](https://docs.spring.io/spring-framework/reference/data-access/transaction.html)
|
||
- [Spring Retry](https://github.com/spring-projects/spring-retry)
|
||
- [ShedLock](https://github.com/lukas-krecan/ShedLock)
|
||
- [Caffeine](https://github.com/ben-manes/caffeine/wiki)
|
||
- [Stripe: Designing robust and predictable APIs with idempotency](https://stripe.com/blog/idempotency)
|