You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 05:36:13