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

Rx框架下并发下载MPEG传输流文件并维持FIFO顺序的实现方案

你遇到的这个问题太典型了——想要用并发提升HLS分片的下载效率,又得保证结果严格按照原始播放列表的顺序递交给订阅者,这在Rx编程里是有标准解法的,不用自己折腾复杂的索引追踪逻辑!

核心实现思路

要实现「并发执行+有序输出」,只需要三步:

  1. 给每个分片绑定原始索引:在序列早期就为每个播放列表项标记它在原始列表中的位置,这样下载完成后能明确它的输出顺序。
  2. 并发执行下载任务:用Merge(maxConcurrency)启动指定数量的并发下载,同时让每个任务返回包含索引和流的结果。
  3. 按索引顺序发射结果:通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 03:53:08