如何实现IEnumerable<T>非破坏性分块?故障时不丢失数据
Chunk运算符在生产者异常时丢失最后一批数据的解决方案
问题背景
在生产者-消费者场景中,生产者是IEnumerable<Item>类型的可枚举序列,使用.NET 6新增的LINQ Chunk运算符按每10个项分块处理。但当生产者中途故障抛出异常时,消费者无法获取故障前生成的最后一批未填满的项——例如生产者生成15个项后抛出异常,消费者仅能收到前10项的分块,11-15项直接丢失。
复现代码
static IEnumerable<int> Produce() { int i = 0; while (true) { i++; Console.WriteLine($"Producing #{i}"); yield return i; if (i == 15) throw new Exception("Oops!"); } } // 消费逻辑 foreach (int[] chunk in Produce().Chunk(10)) { Console.WriteLine($"Consumed: [{String.Join(", ", chunk)}]"); }
实际输出
Producing #1 Producing #2 Producing #3 Producing #4 Producing #5 Producing #6 Producing #7 Producing #8 Producing #9 Producing #10 Consumed: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10] Producing #11 Producing #12 Producing #13 Producing #14 Producing #15 Unhandled exception. System.Exception: Oops! at Program.<Main>g__Produce|0_0()+MoveNext() at System.Linq.Enumerable.ChunkIterator[TSource](IEnumerable`1 source, Int32 size)+MoveNext() at Program.Main()
期望输出
Producing #1 Producing #2 Producing #3 Producing #4 Producing #5 Producing #6 Producing #7 Producing #8 Producing #9 Producing #10 Consumed: [1, 2, 3, 4, 5, 6, 7, 8, 9, 10] Producing #11 Producing #12 Producing #13 Producing #14 Producing #15 Consumed: [11, 12, 13, 14, 15] Unhandled exception. System.Exception: Oops! at Program.<Main>g__Produce|0_0()+MoveNext() at Program.ChunkNonDestructiveIterator[TSource](IEnumerable`1 source, Int32 size)+MoveNext() at Program.Main()
问题
- 能否配置原生
Chunk运算符优先输出缓存的最后一批数据,再抛出异常? - 若不能,如何实现自定义LINQ运算符
ChunkNonDestructive满足需求?
解决方案
原生Chunk的局限性
原生Chunk运算符无法配置此行为,它的迭代逻辑在捕获到生产者异常时会直接抛出,不会先输出已缓存的未填满块。包括System.Interactive的Buffer、MoreLinq的Batch在内的同类运算符,均存在相同的行为。
自定义ChunkNonDestructive实现
我们可以实现一个扩展方法,在迭代序列时捕获异常,先输出已缓存的非空块,再重新抛出异常。代码如下:
using System; using System.Collections.Generic; using System.Linq; public static class EnumerableExtensions { public static IEnumerable<TSource[]> ChunkNonDestructive<TSource>( this IEnumerable<TSource> source, int size) { if (source == null) throw new ArgumentNullException(nameof(source)); if (size < 1) throw new ArgumentOutOfRangeException(nameof(size), "Size must be greater than 0."); return ChunkNonDestructiveIterator(source, size); } private static IEnumerable<TSource[]> ChunkNonDestructiveIterator<TSource>( IEnumerable<TSource> source, int size) { using var enumerator = source.GetEnumerator(); List<TSource> buffer = new List<TSource>(size); while (true) { bool hasNext; try { hasNext = enumerator.MoveNext(); } catch { // 捕获异常时,若缓存有数据则先输出 if (buffer.Count > 0) { yield return buffer.ToArray(); buffer.Clear(); } // 重新抛出原异常 throw; } if (!hasNext) { // 正常结束时输出剩余数据 if (buffer.Count > 0) yield return buffer.ToArray(); yield break; } buffer.Add(enumerator.Current); if (buffer.Count == size) { yield return buffer.ToArray(); buffer.Clear(); } } } }
代码说明
- 参数验证:先检查输入参数的合法性,避免空引用或无效分块大小。
- 迭代逻辑:使用枚举器遍历源序列,将元素添加到缓存列表。
- 异常处理:在调用
MoveNext()时捕获异常,若缓存中有未输出的元素,先输出该块再重新抛出异常。 - 正常结束处理:当序列正常结束时,输出缓存中剩余的所有元素。
测试验证
将消费逻辑改为使用自定义运算符:
foreach (int[] chunk in Produce().ChunkNonDestructive(10)) { Console.WriteLine($"Consumed: [{String.Join(", ", chunk)}]"); }
运行后即可得到期望的输出,先收到11-15的分块,再抛出异常。
内容的提问来源于stack exchange,提问作者Theodor Zoulias
相关产品推荐
相关产品推荐

