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

如何实现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()

问题

  1. 能否配置原生Chunk运算符优先输出缓存的最后一批数据,再抛出异常?
  2. 若不能,如何实现自定义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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 10:45:44