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

C#使用async/await实现UDP对接索尼PTZ相机时匹配序列号等待响应的方法

实现方案

核心逻辑使用ConcurrentDictionary存储每个请求序列号对应的TaskCompletionSource(TCS),通过事件回调触发TCS完成,完全避免轮询,天然适配async/await模式,响应延迟和性能都满足2ms内的要求。

步骤1:新增私有字段

private readonly UDPServer server;
private int _sequenceNum;
// 线程安全字典,key为请求序列号,value为对应等待的任务源
private readonly ConcurrentDictionary<int, TaskCompletionSource<AperatureValue>> _pendingRequests = new();

步骤2:实现异步查询方法

public async Task<AperatureValue> GetCurrentAperatureValueAsync(CancellationToken cancellationToken = default)
{
    // 原子递增生成唯一序列号,避免多线程并发请求时序列号重复
    int currentSeq = Interlocked.Increment(ref _sequenceNum);
    
    // 创建任务源,指定RunContinuationsAsynchronously避免阻塞UDP接收线程
    var tcs = new TaskCompletionSource<AperatureValue>(TaskCreationOptions.RunContinuationsAsynchronously);
    
    // 先把任务源加入待处理字典,再发请求,避免响应太快还没存进去就匹配不到
    if (!_pendingRequests.TryAdd(currentSeq, tcs))
    {
        throw new InvalidOperationException("序列号冲突,请求发送失败");
    }

    try
    {
        // 构造报文,把currentSeq写入报文对应的序列号字段
        byte[] buf = BuildRequestBuffer(currentSeq);
        server.Send("192.168.1.28", 5321, buf);

        // 配置超时逻辑,可根据实际场景调整超时时间
        using var timeoutCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
        timeoutCts.CancelAfter(100);
        
        // 等待响应或者超时/取消
        var completedTask = await Task.WhenAny(tcs.Task, Task.Delay(Timeout.Infinite, timeoutCts.Token));
        if (completedTask != tcs.Task)
        {
            // 超时或者主动取消,清理字典
            _pendingRequests.TryRemove(currentSeq, out _);
            throw new TimeoutException("请求相机超时");
        }

        return await tcs.Task;
    }
    catch
    {
        // 异常情况下清理字典避免内存泄漏
        _pendingRequests.TryRemove(currentSeq, out _);
        throw;
    }
}

// 自行实现构造请求报文的逻辑,把序列号写到对应位置
private byte[] BuildRequestBuffer(int sequenceNum)
{
    // 你的报文构造逻辑
}

步骤3:修改MessageReceived事件处理逻辑

private void Server_MessageReceived(object sender, UDPMessageEventArgs e)
{
    int newSequenceNum = BitConverter.ToInt32(e.sequenceNum);                
    Console.WriteLine("message received" + newSequenceNum + " "+ e.RemoteEndPoint);

    // 匹配待处理请求
    if (_pendingRequests.TryRemove(newSequenceNum, out var tcs))
    {
        try
        {
            // 自行实现把UDP报文解析成AperatureValue的逻辑
            AperatureValue result = ParseResponse(e.Message);
            // 设置任务完成,触发await的地方继续执行
            tcs.TrySetResult(result);
        }
        catch (Exception ex)
        {
            // 解析出错的话设置异常
            tcs.TrySetException(ex);
        }
    }
}

// 自行实现解析响应报文的逻辑
private AperatureValue ParseResponse(byte[] responseData)
{
    // 你的报文解析逻辑
}

注意事项

  • 用Interlocked.Increment保证序列号生成的线程安全,支持并发发起多个查询请求不会冲突
  • 所有异常和超时分支都要清理_pendingRequests字典,避免内存泄漏
  • 如果需要支持其他类型的查询请求,可以把字典的value改成通用类型(比如TaskCompletionSource<object>或者自定义的请求基类),扩展更方便
  • 低概率出现序列号回绕重复的问题,实际使用时如果请求量不大完全可以忽略,要解决的话可以把序列号改成64位长整型,或者重复时做重试逻辑

内容的提问来源于stack exchange,提问作者Chris Kooken

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 14:48:05