如何在Spring Data Redis Reactive中使用Redis事务?
在Spring Boot Reactive Redis中实现事务
ReactiveRedisTemplate确实没有像RedisTemplate那样直接暴露multi()/exec()方法,但你完全可以复用Spring Boot自动配置的连接工厂,基于底层ReactiveRedisConnection实现事务,不需要单独维护额外的Lettuce客户端。
核心实现思路
Redis事务是基于连接的,你可以通过ReactiveRedisTemplate的连接工厂获取ReactiveRedisConnection,调用其multi()开启事务,执行所有需要在事务内的命令,最后用exec()提交或discard()回滚。所有事务内的操作会自动绑定到当前连接的上下文。
基础事务示例
import org.springframework.data.redis.core.ReactiveStringRedisTemplate; import org.springframework.stereotype.Service; import reactor.core.publisher.Mono; @Service public class ReactiveRedisTxService { private final ReactiveStringRedisTemplate reactiveRedisTemplate; public ReactiveRedisTxService(ReactiveStringRedisTemplate reactiveRedisTemplate) { this.reactiveRedisTemplate = reactiveRedisTemplate; } // 执行简单的事务操作 public Mono<Void> runTransactionalSet(String key1, String val1, String key2, String val2) { var connection = reactiveRedisTemplate.getConnectionFactory().getReactiveConnection(); return connection.multi() // 事务内执行两个set操作 .then(reactiveRedisTemplate.opsForValue().set(key1, val1)) .then(reactiveRedisTemplate.opsForValue().set(key2, val2)) // 提交事务 .then(connection.exec()) // 忽略exec返回结果,返回空信号 .then(); } }
带回滚逻辑的示例
如果需要在某些条件下放弃事务,使用discard()方法回滚:
// 带条件回滚的事务 public Mono<Void> runTxWithRollback(String key1, String val1, String key2, String val2) { var connection = reactiveRedisTemplate.getConnectionFactory().getReactiveConnection(); return connection.multi() .then(reactiveRedisTemplate.opsForValue().set(key1, val1)) .flatMap(setResult -> { // 模拟业务校验:如果val2长度不足3则回滚 if (val2.length() < 3) { return connection.discard(); } return reactiveRedisTemplate.opsForValue().set(key2, val2); }) .then(connection.exec()) .then(); }
封装事务模板(可选)
为了避免重复代码,可以封装一个事务模板类,简化事务调用:
import org.springframework.data.redis.core.ReactiveRedisConnection; import org.springframework.data.redis.core.ReactiveRedisTemplate; import reactor.core.publisher.Mono; public class ReactiveRedisTxTemplate { private final ReactiveRedisTemplate<String, String> redisTemplate; public ReactiveRedisTxTemplate(ReactiveRedisTemplate<String, String> redisTemplate) { this.redisTemplate = redisTemplate; } public <T> Mono<T> execute(Mono<T> txOperations) { ReactiveRedisConnection connection = redisTemplate.getConnectionFactory().getReactiveConnection(); return connection.multi() .then(txOperations) .flatMap(result -> connection.exec().thenReturn(result)) // 出错时自动回滚并抛出异常 .onErrorResume(e -> connection.discard().then(Mono.error(e))); } }
使用模板的示例:
@Service public class ReactiveRedisTxService { private final ReactiveStringRedisTemplate reactiveRedisTemplate; private final ReactiveRedisTxTemplate txTemplate; public ReactiveRedisTxService(ReactiveStringRedisTemplate reactiveRedisTemplate) { this.reactiveRedisTemplate = reactiveRedisTemplate; this.txTemplate = new ReactiveRedisTxTemplate(reactiveRedisTemplate); } public Mono<Void> runTxWithTemplate(String key1, String val1, String key2, String val2) { return txTemplate.execute( reactiveRedisTemplate.opsForValue().set(key1, val1) .then(reactiveRedisTemplate.opsForValue().set(key2, val2)) .then() ); } }
注意事项
- Redis事务是队列执行:所有事务内的命令会先进入队列,只有调用
exec()时才会批量执行,中途出错不会自动回滚,需要手动调用discard()放弃队列。 - 确保所有事务操作在同一个Reactor流中执行,保证绑定到同一个Redis连接。
- 完全复用Spring自动配置的连接工厂,不需要额外维护独立的Lettuce客户端。
内容的提问来源于stack exchange,提问作者psyskeptic
相关产品推荐
相关产品推荐

