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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:31:13