基于Redis Stream的事件驱动架构:新消息触发OnMessageReceived事件
基于Redis Stream构建事件驱动架构:新消息触发OnMessageReceived事件
完整实现示例
首先定义自定义事件参数类,用于携带Redis Stream的消息数据:
// 消息接收事件参数,封装Redis Stream消息的ID和内容 public class MqMessageReceivedEventArgs : EventArgs { public string MessageId { get; } public Dictionary<string, RedisValue> MessageData { get; } public MqMessageReceivedEventArgs(string messageId, Dictionary<string, RedisValue> messageData) { MessageId = messageId; MessageData = messageData; } }
核心Redis Stream监听类实现:
using StackExchange.Redis; using System; using System.Threading; using System.Threading.Tasks; public class RedisStreamEventListener { private readonly IDatabase _redisDatabase; private readonly string _streamName; private string _lastProcessedId = "$"; // "$"代表从当前最新消息开始监听 // 定义消息接收事件 public event EventHandler<MqMessageReceivedEventArgs> OnMessageReceived; public RedisStreamEventListener(IConnectionMultiplexer redisConnection, string streamName) { _redisDatabase = redisConnection.GetDatabase(); _streamName = streamName; } // 启动持续监听的方法 public async Task StartListeningAsync(CancellationToken cancellationToken = default) { while (!cancellationToken.IsCancellationRequested) { try { // 阻塞读取新消息(超时5秒,无消息时等待,避免无效轮询) var messages = await _redisDatabase.StreamReadAsync(_streamName, _lastProcessedId, count: 100, timeout: 5000); foreach (var message in messages) { // 更新最后处理的消息ID,确保下次从最新位置开始消费 _lastProcessedId = message.Id; // 触发消息接收事件 OnMessageReceivedEvent(new MqMessageReceivedEventArgs(message.Id, message.ToDictionary())); } } catch (Exception ex) { // 可替换为日志记录等异常处理逻辑 Console.WriteLine($"Redis Stream监听异常: {ex.Message}"); await Task.Delay(1000, cancellationToken); // 异常后延迟重试 } } } // 受保护虚方法,允许子类重写事件触发逻辑 protected virtual void OnMessageReceivedEvent(MqMessageReceivedEventArgs e) { OnMessageReceived?.Invoke(this, e); } }
关键说明
- 实时监听的正确方式:你提供的
StreamRangeAsync用于查询历史消息,无法实现实时监听。推荐使用StreamReadAsync带超时的阻塞读取,有新消息时立即返回,无消息则等待超时,避免无效轮询浪费资源。 - 消息位置跟踪:用
$初始化_lastProcessedId表示从当前最新消息开始监听;每次处理完消息后更新该ID,确保不会重复消费消息。 - 事件设计规范:
OnMessageReceivedEvent虚方法的设计允许子类重写事件触发逻辑,同时保留事件的多播委托特性,确保线程安全。 - 异常容错:监听循环中加入异常捕获与重试逻辑,避免单次异常导致监听中断。
使用示例
// 初始化Redis连接 var redisConnection = await ConnectionMultiplexer.ConnectAsync("localhost:6379"); var listener = new RedisStreamEventListener(redisConnection, "your-target-stream"); // 订阅消息接收事件 listener.OnMessageReceived += (sender, e) => { Console.WriteLine($"收到消息ID: {e.MessageId}"); foreach (var kvp in e.MessageData) { Console.WriteLine($"{kvp.Key}: {kvp.Value}"); } }; // 启动监听(可结合后台服务或HostedService在.NET应用中持久运行) await listener.StartListeningAsync();
内容的提问来源于stack exchange,提问作者Mahesh
相关产品推荐
相关产品推荐

