RestSharp处理分块JSON事件流:实时读取与多对象解析
使用RestSharp处理分块HTTP事件流:实时读取与类型化解析
我来帮你搞定这两个RestSharp处理分块事件流的问题,这在处理异步分块API或者类SSE场景时特别常见,咱们一步步拆解解决:
1. 实时读取分块数据(不等待连接关闭)
RestSharp默认的ExecuteAsync会等整个响应完成才返回结果,但对于Transfer-Encoding: chunked的场景,我们需要直接操作响应流来实现实时读取。
你可以用RestClient.ExecuteStreamAsync获取响应流,然后循环读取流内容。需要注意HTTP分块的格式:每个块以十六进制长度开头,后跟\r\n,接着是块内容,最后再以\r\n结尾。不过从你提供的示例来看,服务器返回的内容已经简化了展示,你可以根据实际情况选择读取方式:
方式1:处理标准HTTP分块格式
如果服务器严格遵循分块规范,用以下逻辑读取:
var client = new RestClient("https://myserver.com:8014"); var request = new RestRequest("/OTEvents", Method.Get); request.AddQueryParameter("subscriptionId", "eyJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJhZG1pbiIsImlhdCI6MTUyNzY2MTg0MX0.fxsP4bLNzqSSFtYsTNmyyV4bM-OBwhcwhy-w_HwQYmQ"); using var response = await client.ExecuteStreamAsync(request); using var reader = new StreamReader(response); while (!reader.EndOfStream) { // 读取分块长度(十六进制) var lengthLine = await reader.ReadLineAsync(); if (string.IsNullOrWhiteSpace(lengthLine)) continue; if (!int.TryParse(lengthLine, System.Globalization.NumberStyles.HexNumber, null, out var chunkLength)) continue; // 跳过无效分块 // 读取对应长度的JSON内容 var buffer = new char[chunkLength]; await reader.ReadAsync(buffer, 0, chunkLength); var jsonContent = new string(buffer); // 跳过分块结尾的\r\n await reader.ReadLineAsync(); // 处理单个JSON块 ProcessEventJson(jsonContent); }
方式2:简化逐行读取(如果JSON块是单独行)
如果服务器返回的分块内容已经去掉长度前缀,且每个JSON块单独占一行,直接用ReadLineAsync更简单:
using var response = await client.ExecuteStreamAsync(request); using var reader = new StreamReader(response); string line; while ((line = await reader.ReadLineAsync()) != null) { if (string.IsNullOrWhiteSpace(line)) continue; ProcessEventJson(line); }
2. 根据事件类型解析为不同对象并触发回调
首先为每个事件类型定义对应实体类,再通过事件名称动态解析,最后触发回调:
步骤1:定义事件实体类
// 所有事件的基类,用于获取事件名称 public class EventBase { public string eventName { get; set; } } // 具体事件类,根据你的JSON结构补充属性 public class OnChannelInformationEvent : EventBase { public string text { get; set; } } public class OnCallCreatedEvent : EventBase { public string loginName { get; set; } public string callRef { get; set; } public CallData callData { get; set; } public List<Leg> legs { get; set; } public List<Participant> participants { get; set; } } public class OnCallModifiedEvent : EventBase { public string loginName { get; set; } public string callRef { get; set; } public List<Leg> addedLegs { get; set; } } // 辅助实体类,根据JSON结构定义 public class CallData { public InitialCalled initialCalled { get; set; } public string state { get; set; } } public class InitialCalled { public PhoneId id { get; set; } } public class PhoneId { public string phoneNumber { get; set; } } public class Leg { public string deviceId { get; set; } public string media { get; set; } } public class Participant { public string participantId { get; set; } public Identity identity { get; set; } } public class Identity { public PhoneId id { get; set; } public string firstName { get; set; } }
步骤2:解析JSON并触发回调
这里以Newtonsoft.Json为例(也可以用System.Text.Json,逻辑类似):
// 定义回调委托 public delegate void EventReceivedHandler(object eventObj); // 暴露回调事件 public event EventReceivedHandler OnEventReceived; private void ProcessEventJson(string json) { try { // 先反序列化为基类,获取事件名称 var baseEvent = JsonConvert.DeserializeObject<EventBase>(json); if (baseEvent == null) return; // 根据事件名称解析为具体类 object specificEvent = baseEvent.eventName switch { "OnChannelInformation" => JsonConvert.DeserializeObject<OnChannelInformationEvent>(json), "OnCallCreated" => JsonConvert.DeserializeObject<OnCallCreatedEvent>(json), "OnCallModified" => JsonConvert.DeserializeObject<OnCallModifiedEvent>(json), _ => baseEvent // 未知事件返回基类 }; // 触发回调通知主逻辑 OnEventReceived?.Invoke(specificEvent); } catch (JsonException ex) { Console.WriteLine($"解析事件失败: {ex.Message}"); } }
步骤3:订阅回调处理事件
// 在主逻辑中订阅回调 OnEventReceived += (eventObj) => { switch (eventObj) { case OnChannelInformationEvent infoEvent: Console.WriteLine($"收到通道事件: {infoEvent.text}"); break; case OnCallCreatedEvent createdEvent: Console.WriteLine($"呼叫创建ID: {createdEvent.callRef}, 状态: {createdEvent.callData.state}"); break; case OnCallModifiedEvent modifiedEvent: Console.WriteLine($"呼叫修改ID: {modifiedEvent.callRef}, 新增设备: {modifiedEvent.addedLegs[0].deviceId}"); break; default: Console.WriteLine($"未知事件: {((EventBase)eventObj).eventName}"); break; } };
完整整合示例
把上面的代码整合起来,就是一个可运行的完整流程:
using RestSharp; using Newtonsoft.Json; using System; using System.Collections.Generic; using System.IO; using System.Threading.Tasks; public class EventStreamProcessor { public delegate void EventReceivedHandler(object eventObj); public event EventReceivedHandler OnEventReceived; public async Task StartListeningAsync(string subscriptionId) { var client = new RestClient("https://myserver.com:8014"); var request = new RestRequest("/OTEvents", Method.Get); request.AddQueryParameter("subscriptionId", subscriptionId); try { using var response = await client.ExecuteStreamAsync(request); using var reader = new StreamReader(response); while (!reader.EndOfStream) { var lengthLine = await reader.ReadLineAsync(); if (string.IsNullOrWhiteSpace(lengthLine)) continue; if (!int.TryParse(lengthLine, System.Globalization.NumberStyles.HexNumber, null, out var chunkLength)) continue; var buffer = new char[chunkLength]; await reader.ReadAsync(buffer, 0, chunkLength); var jsonContent = new string(buffer); await reader.ReadLineAsync(); // 跳过结尾的\r\n ProcessEventJson(jsonContent); } } catch (Exception ex) { Console.WriteLine($"连接异常: {ex.Message}"); } } private void ProcessEventJson(string json) { try { var baseEvent = JsonConvert.DeserializeObject<EventBase>(json); if (baseEvent == null) return; object specificEvent = baseEvent.eventName switch { "OnChannelInformation" => JsonConvert.DeserializeObject<OnChannelInformationEvent>(json), "OnCallCreated" => JsonConvert.DeserializeObject<OnCallCreatedEvent>(json), "OnCallModified" => JsonConvert.DeserializeObject<OnCallModifiedEvent>(json), _ => baseEvent }; OnEventReceived?.Invoke(specificEvent); } catch (JsonException ex) { Console.WriteLine($"JSON解析失败: {ex.Message}"); } } // 事件实体类定义 public class EventBase { public string eventName { get; set; } } public class OnChannelInformationEvent : EventBase { public string text { get; set; } } public class OnCallCreatedEvent : EventBase { public string loginName { get; set; } public string callRef { get; set; } public CallData callData { get; set; } public List<Leg> legs { get; set; } public List<Participant> participants { get; set; } } public class OnCallModifiedEvent : EventBase { public string loginName { get; set; } public string callRef { get; set; } public List<Leg> addedLegs { get; set; } } public class CallData { public InitialCalled initialCalled { get; set; } public string state { get; set; } } public class InitialCalled { public PhoneId id { get; set; } } public class PhoneId { public string phoneNumber { get; set; } } public class Leg { public string deviceId { get; set; } public string media { get; set; } } public class Participant { public string participantId { get; set; } public Identity identity { get; set; } } public class Identity { public PhoneId id { get; set; } public string firstName { get; set; } } } // 使用示例 class Program { static async Task Main(string[] args) { var processor = new EventStreamProcessor(); processor.OnEventReceived += (eventObj) => { switch (eventObj) { case EventStreamProcessor.OnChannelInformationEvent info: Console.WriteLine($"[通道信息] {info.text}"); break; case EventStreamProcessor.OnCallCreatedEvent created: Console.WriteLine($"[呼叫创建] ID: {created.callRef}, 状态: {created.callData.state}"); break; case EventStreamProcessor.OnCallModifiedEvent modified: Console.WriteLine($"[呼叫修改] ID: {modified.callRef}, 新增设备: {modified.addedLegs[0].deviceId}"); break; default: Console.WriteLine($"[未知事件] {((EventStreamProcessor.EventBase)eventObj).eventName}"); break; } }; var subscriptionId = "eyJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJhZG1pbiIsImlhdCI6MTUyNzY2MTg0MX0.fxsP4bLNzqSSFtYsTNmyyV4bM-OBwhcwhy-w_HwQYmQ"; await processor.StartListeningAsync(subscriptionId); } }
注意事项
- 确保使用的RestSharp版本支持
ExecuteStreamAsync(v107+版本都支持)。 - 如果需要断线重连,可以在异常捕获块中添加重试逻辑。
- 若使用System.Text.Json,只需将
JsonConvert替换为JsonSerializer即可,核心逻辑一致。
内容的提问来源于stack exchange,提问作者Hans-Jürgen Brand
相关产品推荐
相关产品推荐

