如何从Redis背板反序列化SignalR消息
解析SignalR Redis背板中的MessagePack消息
一、消息格式说明
你看到的开头\x92\x90\x81\xa4json和结尾字节是MessagePack序列化格式的标识,SignalR的Redis背板默认用这种格式包装消息:
\x92:表示当前是包含2个元素的数组\x90:数组第一个元素是空数组(SignalR内部元数据,此处无实际内容)\x81:数组第二个元素是包含1个键值对的Map\xa4json:Map的键是长度为4的字符串"json",对应的值就是你需要的实际消息JSON- 末尾的
\x1e是Redis消息的结束符,反序列化时会自动忽略
二、反序列化解决方案
1. 定义消息实体类
根据示例JSON格式,创建对应的实体类:
using System.Text.Json.Serialization; public class Message { [JsonPropertyName("type")] public int Type { get; set; } [JsonPropertyName("target")] public string Target { get; set; } [JsonPropertyName("arguments")] public List<MessageArgument> Arguments { get; set; } } public class MessageArgument { [JsonPropertyName("senderId")] public string SenderId { get; set; } [JsonPropertyName("message")] public string MessageContent { get; set; } [JsonPropertyName("sentAt")] public DateTime SentAt { get; set; } }
2. 实现消息解析方法
修改ReadMessage方法,先解析MessagePack数据,提取JSON后再转为实体:
using MessagePack; using System.Text.Json; private static Message ReadMessage(RedisValue value) { // 将RedisValue转为字节数组 byte[] messagePackBytes = value; // 反序列化为SignalR固定的MessagePack数组结构:[元数据数组, { "json": 实际消息JSON }] var messagePackArray = MessagePackSerializer.Deserialize<object[]>(messagePackBytes); // 提取包含json键的字典 var jsonContainer = messagePackArray[1] as Dictionary<string, object>; if (jsonContainer == null || !jsonContainer.TryGetValue("json", out var jsonObj)) throw new InvalidOperationException("无法提取SignalR消息中的JSON内容"); // 将JSON字符串反序列化为Message实体 string jsonString = jsonObj.ToString(); return JsonSerializer.Deserialize<Message>(jsonString, new JsonSerializerOptions { PropertyNameCaseInsensitive = true }); }
3. 修正订阅回调的问题
原代码中static async void Handler存在异常无法捕获、无法访问外部connection变量的问题,改为异步实例方法:
protected override async Task ExecuteAsync(CancellationToken stoppingToken) { using (var connection = new NpgsqlConnection(_configuration.GetConnectionString("YugabyteDB"))) { await connection.OpenAsync(stoppingToken); // 异步回调方法,可访问外部connection变量 async Task MessageHandler(RedisChannel channel, RedisValue value) { if (value == default) return; try { var message = ReadMessage(value); // 替换为你的实际插入SQL和参数 const string InsertSql = @" INSERT INTO messages (target, sender_id, message_content, sent_at) VALUES (@Target, @SenderId, @MessageContent, @SentAt)"; await connection.ExecuteAsync(InsertSql, new { message.Target, message.Arguments[0].SenderId, message.Arguments[0].MessageContent, message.Arguments[0].SentAt }, cancellationToken: stoppingToken); } catch (Exception ex) { // 处理解析或插入异常,比如日志记录 Console.WriteLine($"处理SignalR消息失败:{ex.Message}"); } } var subscriber = _connectionMultiplexer.GetSubscriber(); // 订阅用户消息频道 await subscriber.SubscribeAsync( new RedisChannel("*SignalRHub.MessagingHub:user:*", RedisChannel.PatternMode.Pattern), MessageHandler, stoppingToken); // 订阅群组消息频道 await subscriber.SubscribeAsync( new RedisChannel("*SignalRHub.MessagingHub:group:*", RedisChannel.PatternMode.Pattern), MessageHandler, stoppingToken); await Task.Delay(Timeout.Infinite, stoppingToken); } }
三、必要的NuGet包
确保项目安装以下依赖:
MessagePack:用于反序列化MessagePack数据System.Text.Json:用于解析JSON消息(若使用Newtonsoft.Json则安装Newtonsoft.Json)StackExchange.Redis:Redis连接库Dapper:若使用ExecuteAsync执行SQL(或替换为EF Core等其他ORM)
内容的提问来源于stack exchange,提问作者Szyszka947
相关产品推荐
相关产品推荐

