TPL DataFlow处理HTTP加载消息不全问题求助
问题根源分析
你遇到的随机丢数据问题,核心原因是你只等待了上游downloader块的完成,却没等整个数据流管道完全处理完所有数据。
当你调用downloader.Complete()后,downloader会停止接收新请求,处理完剩余下载任务后就进入完成状态。但此时deserializer可能还在处理最后几个下载回来的字符串,甚至还没把反序列化好的DummyData发送到buffer里。这时候直接调用buffer.TryReceiveAll(),只能拿到已经进入buffer的数据,那些还在deserializer里处理的结果就被漏掉了——这就是丢失数量随机的原因,完全看最后几个任务的处理速度。
修复方案
最直接的修复就是等待最后一个块(也就是buffer)的完成,因为你已经设置了PropagateCompletion = true,完成信号会从downloader传递到deserializer,再传递到buffer。只有当buffer.Completion完成时,才能确保所有数据都已经处理完毕并进入buffer。
修改后的代码如下:
public async Task<IEnumerable<DummyData>> LoadAsync(IEnumerable<Uri> uris) { IList<DummyData> result; using (var client = new HttpClient()) { var buffer = new BufferBlock<DummyData>(); var downloader = new TransformBlock<Uri, string>( async u => await client.GetStringAsync(u), new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = _maxParallelism }); var deserializer = new TransformBlock<string, DummyData>( s => JsonConvert.DeserializeObject<DummyData>(s), new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = DataflowBlockOptions.Unbounded }); var linkOptions = new DataflowLinkOptions { PropagateCompletion = true }; downloader.LinkTo(deserializer, linkOptions); deserializer.LinkTo(buffer, linkOptions); foreach (Uri uri in uris) { await downloader.SendAsync(uri); } downloader.Complete(); // 等待最后一个块完成,确保所有数据都已处理并进入buffer await buffer.Completion; buffer.TryReceiveAll(out result); } return result; }
额外优化建议
- 替换BufferBlock为直接收集结果:其实你可以不用
BufferBlock,改用ActionBlock直接把结果添加到线程安全列表里,这样更直观,也避免TryReceiveAll的潜在问题:
public async Task<IEnumerable<DummyData>> LoadAsync(IEnumerable<Uri> uris) { var resultList = new ConcurrentBag<DummyData>(); using (var client = new HttpClient()) { var downloader = new TransformBlock<Uri, string>( async u => await client.GetStringAsync(u), new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = _maxParallelism }); var deserializer = new TransformBlock<string, DummyData>( s => JsonConvert.DeserializeObject<DummyData>(s), new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = DataflowBlockOptions.Unbounded }); var collector = new ActionBlock<DummyData>( data => resultList.Add(data), new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = DataflowBlockOptions.Unbounded }); var linkOptions = new DataflowLinkOptions { PropagateCompletion = true }; downloader.LinkTo(deserializer, linkOptions); deserializer.LinkTo(collector, linkOptions); foreach (Uri uri in uris) { await downloader.SendAsync(uri); } downloader.Complete(); await collector.Completion; } return resultList; }
- 添加错误处理:你当前的代码没有任何错误处理,如果某个下载请求失败或者反序列化抛出异常,整个数据流管道会进入故障状态,后续任务都会被丢弃,这也可能导致数据丢失。建议在块的委托里添加
try-catch,或者设置块的故障处理逻辑。
内容的提问来源于stack exchange,提问作者Zdeněk
相关产品推荐
相关产品推荐

