Redis 7.2 Stream 消息队列实战:替代 Kafka 的轻量方案

Redis 7.2 Streams as Message Queue: A Lightweight Alternative to Kafka

| Alex Chen | 2026-08-03T09:34:38

不是所有场景都需要 Kafka。这篇分享我们用 Redis Stream 实现消息队列的完整实践,以及什么时候该用什么时候不该用。

Not every scenario needs Kafka. A complete guide to implementing message queues with Redis Streams and when to use them.

## 背景 我们有个内部通知系统,需求很简单: - 接收各个服务的通知事件 - 按用户分发推送(邮件、站内信、App Push) - 日消息量约 10 万条 - 允许偶尔丢消息(不是支付场景) 之前用的 Kafka,三节点集群,ZooKeeper 单独部署。维护成本高得离谱,光是 Kafka 集群就占了 12GB 内存。对于日均 10 万条消息来说,像开着卡车送外卖。 ## 为什么选 Redis Stream Redis Stream 是 Redis 5.0 引入的数据结构,7.2 版本进一步增强。它的核心能力: - **持久化消息日志**:类似 Kafka 的 append-only log - **消费者组**:支持多个消费者分摊消息处理 - **ACK 机制**:消息确认,处理失败可重试 - **历史回溯**:可以从任意位置重新消费 而且它就是 Redis 的一个数据结构,不需要额外部署任何东西。 ## 生产者端 ```java @Service @RequiredArgsConstructor public class NotificationProducer { private final StringRedisTemplate redis; private static final String STREAM_KEY = "notification:stream"; public String send(NotificationEvent event) { Map message = Map.of( "userId", event.getUserId().toString(), "type", event.getType().name(), "title", event.getTitle(), "content", event.getContent(), "timestamp", String.valueOf(System.currentTimeMillis()) ); // XADD 追加消息,MAXLEN ~ 100000 限制流长度(自动淘汰旧消息) RecordId id = redis.opsForStream().add( StreamRecords.mapBacked(message).withStreamKey(STREAM_KEY) ); return id.getValue(); } } ``` `MAXLEN ~ 100000` 里的 `~` 表示近似裁剪,Redis 会在内部用更高效的方式清理,不会每条消息都精确裁剪。 ## 消费者端 ```java @Component @Slf4j public class NotificationConsumer implements ApplicationRunner { private final StringRedisTemplate redis; private final NotificationService notificationService; private static final String STREAM = "notification:stream"; private static final String GROUP = "notification-group"; private static final String CONSUMER = "consumer-" + UUID.randomUUID().toString().substring(0, 8); @Override public void run(ApplicationArguments args) { // 创建消费者组(如果不存在) try { redis.opsForStream().createGroup(STREAM, ReadOffset.from("0"), GROUP); } catch (Exception e) { // 组已存在,忽略 } // 启动消费循环 new Thread(this::consumeLoop, "stream-consumer").start(); } private void consumeLoop() { while (!Thread.currentThread().isInterrupted()) { try { // XREADGROUP 阻塞读取,超时 2 秒 List> records = redis.opsForStream().read( Consumer.from(GROUP, CONSUMER), StreamReadOptions.empty().count(10).block(Duration.ofSeconds(2)), StreamOffset.create(STREAM, ReadOffset.lastConsumed()) ); if (records != null) { for (var record : records) { try { processMessage(record.getValue()); // ACK 确认 redis.opsForStream().acknowledge(STREAM, GROUP, record.getId()); } catch (Exception e) { log.error("处理消息失败: {}", record.getId(), e); // 不 ACK,消息会进入 pending 列表,后续可重试 } } } } catch (Exception e) { log.error("消费循环异常", e); sleep(1000); } } } } ``` ## 处理失败消息(Pending 重试) ```java @Scheduled(fixedDelay = 60000) // 每分钟检查一次 public void retryPending() { // 查找超过 30 秒未 ACK 的消息 PendingMessages pending = redis.opsForStream().pending( STREAM, GROUP, Range.closed("-", "+"), 100 ); for (PendingMessage pm : pending) { if (pm.getElapsedTimeSinceLastDelivery().compareTo(Duration.ofSeconds(30)) > 0) { if (pm.getTotalDeliveryCount() >= 3) { // 重试 3 次还失败,移到死信队列 redis.opsForStream().acknowledge(STREAM, GROUP, pm.getId()); moveToDeadLetter(pm); } else { // XCLAIM 认领消息,重新处理 redis.opsForStream().claim(STREAM, GROUP, CONSUMER, Duration.ofSeconds(30), pm.getId()); } } } } ``` ## 监控 ```java @GetMapping("/api/admin/stream/stats") public Result streamStats() { StreamInfo.XInfoStream info = redis.opsForStream().info(STREAM); StreamInfo.XInfoGroups groups = redis.opsForStream().groups(STREAM); Map stats = new HashMap(); stats.put("length", info.streamLength()); stats.put("firstEntry", info.firstEntry()); stats.put("lastEntry", info.lastEntry()); for (var group : groups) { stats.put("group_" + group.groupName() + "_pending", group.pendingCount()); stats.put("group_" + group.groupName() + "_consumers", group.consumerCount()); } return Result.ok(stats); } ``` ## Redis Stream vs Kafka | 维度 | Redis Stream | Kafka | |------|-------------|-------| | 部署复杂度 | 零(用现有 Redis) | 高(ZK/KRaft + Broker) | | 资源占用 | 极低 | 高(JVM) | | 消息持久化 | AOF/RDB | 磁盘日志 | | 吞吐量 | ~10 万/秒 | ~100 万/秒 | | 消息回溯 | 支持 | 更强 | | 消费者组 | 支持 | 更成熟 | | 适合场景 | 中小规模 | 大规模 | ## 什么时候不该用 Redis Stream 1. **日消息量超过百万级**:Redis 是内存数据库,消息量大了内存扛不住 2. **消息不能丢**:Redis 的持久化不如 Kafka 可靠,宕机可能丢几秒数据 3. **需要精确一次语义**:Redis Stream 只保证 at-least-once 4. **需要长期保留消息**:Kafka 可以保留几天甚至更久,Redis 受内存限制 ## 迁移效果 - 内存占用:12GB(Kafka 集群) → 200MB(Redis Stream 数据) - 运维复杂度:3 节点 Kafka + ZK → 无额外组件 - 消息延迟:基本一致,都在毫秒级 选对工具比选好工具更重要。


## Background Our internal notification system handles ~100K messages/day. Running a 3-node Kafka cluster (12GB RAM) for this volume was massive overkill. ## Redis Stream Solution Redis Streams provide append-only log with consumer groups, ACK mechanism, and message replay - all built into existing Redis, zero additional deployment. ## Key Implementation Producer uses XADD with approximate MAXLEN trimming. Consumer uses XREADGROUP with blocking reads and manual ACK. Failed messages enter pending list with automatic retry via XCLAIM and dead letter queue after 3 attempts. ## When NOT to Use Daily volume exceeding millions, zero message loss requirement, exactly-once semantics needed, or long-term message retention required. ## Results Memory: 12GB → 200MB. Operational complexity: eliminated Kafka/ZK entirely. Latency: unchanged at millisecond level.

← Back to News