Rx框架下并发下载MPEG传输流文件并维持FIFO顺序的实现方案
你遇到的这个问题太典型了——想要用并发提升HLS分片的下载效率,又得保证结果严格按照原始播放列表的顺序递交给订阅者,这在Rx编程里是有标准解法的,不用自己折腾复杂的索引追踪逻辑!
核心实现思路
要实现「并发执行+有序输出」,只需要三步:
- 给每个分片绑定原始索引:在序列早期就为每个播放列表项标记它在原始列表中的位置,这样下载完成后能明确它的输出顺序。
- 并发执行下载任务:用
Merge(maxConcurrency)启动指定数量的并发下载,同时让每个任务返回包含索引和流的结果。 - 按索引顺序发射结果:通过
Scan维护一个等待队列和下一个待发射的索引,每当有新结果进入队列,就检查是否能按顺序输出当前等待的结果,确保订阅者收到的顺序和原始列表完全一致。
修改后的完整代码
public async Task<IObservable<(int Current, int Total, Stream Stream)>> FetchVideoSegments() { var playlist = await GetPlaylist(); var total = playlist.PlaylistItems.Count(); return playlist.PlaylistItems.ToObservable() // 1. 为每个分片绑定原始索引(从0开始计数) .Select((item, originalIndex) => (item.Uri, originalIndex)) // 转换为绝对URI .Select(tuple => (AbsoluteUri: MakeRelativeAbsoluteUrl(tuple.Uri), tuple.originalIndex)) // 包装为延迟执行的下载任务,返回带原始索引的流 .Select(tuple => Observable.Defer(() => DownloadVideoSegment(tuple.AbsoluteUri) .Select(stream => (Index: tuple.originalIndex, Stream: stream)) )) // 2. 启动并发下载,限制最大并发数(比如你说的2-3) .Merge(FetchSegmentsMaxConcurrency) // 3. 按原始索引顺序输出结果 .Scan( // 初始状态:下一个要输出的索引是0,等待队列为有序字典(按索引排序) (NextExpectedIndex: 0, PendingResults: new SortedDictionary<int, Stream>()), (currentState, newResult) => { // 将新完成的下载结果加入等待队列 currentState.PendingResults.Add(newResult.Index, newResult.Stream); return currentState; } ) // 从状态中提取可以按顺序输出的结果 .SelectMany(state => { var readyResults = new List<(int Index, Stream Stream)>(); // 检查队列中是否有当前期望的索引,有则取出并输出,直到没有匹配项为止 while (state.PendingResults.ContainsKey(state.NextExpectedIndex)) { var stream = state.PendingResults[state.NextExpectedIndex]; readyResults.Add((state.NextExpectedIndex, stream)); state.PendingResults.Remove(state.NextExpectedIndex); state.NextExpectedIndex++; } return readyResults; }) // 转换为订阅者需要的进度元组(Current从1开始计数,符合用户习惯) .Select(result => (Current: result.Index + 1, Total: total, result.Stream)); }
关键细节解释
- 绑定原始索引:
Select((item, originalIndex) => ...)里的originalIndex是Rx自动提供的序列索引,完全对应播放列表的原始顺序,这是后续有序输出的基础。 - 并发控制:
Merge(FetchSegmentsMaxConcurrency)会同时启动最多指定数量的下载任务,既提升了效率,又不会因为并发过高导致服务器压力过大。 - 有序输出逻辑:
Scan维护的状态里,NextExpectedIndex标记了下一个要递交给订阅者的分片索引,PendingResults则存储已经完成但还没到输出时机的结果(用SortedDictionary保证索引有序)。每当有新结果加入,我们就检查是否能输出当前期望的索引对应的结果,如果可以,就把它从队列中取出并交给订阅者,同时更新期望索引,直到队列中没有匹配的结果为止。
这个方案既保证了并发下载的效率,又严格遵守了原始播放列表的顺序,完全符合你的需求~
内容的提问来源于stack exchange,提问作者Morgan Touverey Quilling
相关产品推荐
相关产品推荐

