Session 共享与消息队列
约 738 字大约 2 分钟
布欧-Lewyon
2026-05-15
Session 共享
微服务多实例部署时,用户可能请求到不同实例,Session 需要共享。
Spring Session Redis
<dependency>
<groupId>org.springframework.session</groupId>
<artifactId>spring-session-data-redis</artifactId>
</dependency>spring:
session:
store-type: redis
redis:
namespace: spring:session
flush-mode: on_save
data:
redis:
host: localhost
port: 6379@Configuration
@EnableRedisHttpSession(maxInactiveIntervalInSeconds = 1800) // 30 分钟
public class SessionConfig {
// 无需额外代码,Session 自动存储到 Redis
}
@RestController
public class LoginController {
@PostMapping("/login")
public String login(HttpSession session, @RequestParam String username) {
session.setAttribute("user", username);
return "登录成功";
}
@GetMapping("/profile")
public String profile(HttpSession session) {
String user = (String) session.getAttribute("user");
return "当前用户: " + user;
}
@PostMapping("/logout")
public String logout(HttpSession session) {
session.invalidate();
return "已退出";
}
}自定义 Session 存储
@Service
public class CustomSessionService {
@Autowired
private StringRedisTemplate redis;
// 创建 Session
public String createSession(String userId, Map<String, Object> data) {
String sessionId = UUID.randomUUID().toString();
String key = "session:" + sessionId;
redis.opsForHash().putAll(key, data);
redis.expire(key, 30, TimeUnit.MINUTES);
return sessionId;
}
// 获取 Session
public Map<Object, Object> getSession(String sessionId) {
String key = "session:" + sessionId;
Map<Object, Object> data = redis.opsForHash().entries(key);
if (!data.isEmpty()) {
redis.expire(key, 30, TimeUnit.MINUTES); // 续期
}
return data;
}
// 删除 Session
public void removeSession(String sessionId) {
redis.delete("session:" + sessionId);
}
}消息队列
List 队列(简单场景)
@Service
public class SimpleQueueService {
@Autowired
private StringRedisTemplate redis;
// 发送任务
public void sendTask(String queueName, String task) {
redis.opsForList().leftPush(queueName, task);
}
// 接收任务(阻塞)
public String receiveTask(String queueName, long timeout) throws InterruptedException {
// BLPOP:阻塞直到有消息
CountDownLatch latch = new CountDownLatch(1);
final String[] result = {null};
// 使用 RedisTemplate 的 execute 调用 BLPOP
// 或用 RedisMessageListenerContainer
return result[0];
}
// 定时轮询
@Scheduled(fixedDelay = 1000)
public void pollTask() {
String task = redis.opsForList().rightPop("task:queue");
if (task != null) {
processTask(task);
}
}
}Stream 队列(可靠消费)
@Service
public class StreamQueueService {
@Autowired
private StringRedisTemplate redis;
private static final String STREAM = "stream:orders";
private static final String GROUP = "order-processors";
@PostConstruct
public void init() {
// 创建消费者组
try {
redis.opsForStream().createGroup(STREAM, GROUP);
} catch (Exception e) {
// 组已存在则忽略
}
}
// 生产者
public void sendOrder(Order order) {
Map<String, String> message = Map.of(
"orderId", order.getId().toString(),
"userId", order.getUserId().toString(),
"amount", order.getAmount().toString()
);
RecordId id = redis.opsForStream()
.add(STREAM, message);
}
// 消费者(定时拉取)
@Scheduled(fixedDelay = 200)
public void consume() {
// 读取未确认消息
List<MapRecord<String, String, String>> messages = redis.opsForStream()
.read(Consumer.from(GROUP, "consumer-1"),
StreamReadOptions.empty().count(10),
StreamOffset.create(STREAM, ReadOffset.lastConsumed()));
for (MapRecord<String, String, String> msg : messages) {
try {
processOrder(msg.getValue());
// ACK 确认
redis.opsForStream().acknowledge(STREAM, GROUP, msg.getId());
} catch (Exception e) {
log.error("处理失败: {}", msg.getId(), e);
// 失败消息留在 Pending 列表,后续重试
}
}
}
// 处理死信(重试多次失败的消息)
@Scheduled(fixedDelay = 60000)
public void retryDeadLetters() {
PendingMessages pending = redis.opsForStream()
.pending(STREAM, GROUP, PendingMessagesSummary.class)
.getPendingMessages();
// 查看长时间未确认的消息并处理
}
}小结
| 功能 | 实现方式 | 适用场景 |
|---|---|---|
| Session 共享 | spring-session-data-redis | 分布式系统 Session |
| 简单队列 | LPUSH + BRPOP | 简单任务分发 |
| 可靠队列 | Stream + 消费者组 + ACK | 需要消息确认的生产级场景 |
| 延迟队列 | SortedSet(时间戳作分数) | 延时任务 |
- Spring Session 零代码实现 Session 共享。
- Stream 队列支持消费者组、ACK、消息回溯,替代 List 队列。
- 延迟队列用 SortedSet 分数存时间戳,轮询到期消息。
上一节:分布式锁与限流器
