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
相关产品推荐
相关产品推荐

