文档智能处理流水线设计

从摄取到输出的端到端Pipeline架构与.NET实战

为什么需要流水线架构?

在文档智能处理场景中,一份PDF文件从进入系统到产生结构化结果,往往要经历十几个处理环节:解压、预处理、页面渲染、OCR识别、文本清洗、字段提取、分类、校验、输出存储等。如果将所有步骤串行硬编码在一个方法中,代码难以维护、无法并行化、出错时必须全链路重试。流水线(Pipeline)架构通过将处理分解为一系列独立、可组合的阶段(Stage),每个阶段只关心单一职责,既提升了可维护性,也释放了并发能力。

本文系统讲解如何在.NET中设计一个生产级的文档处理流水线,涵盖架构设计、阶段划分、异步模式、错误处理、重试与监控,并给出DocCore SDK中实际使用的代码样例。

一、流水线的标准阶段

一个典型的文档智能处理流水线包含以下五个核心阶段:

阶段 主要任务 输入 输出
Ingest(摄取)接入文件、检测类型、基础元数据文件路径/流原始文档对象
Preprocess(预处理)页面渲染、图像增强、倾斜校正原始文档规范化页面
Extract(提取)OCR、表格识别、布局分析规范化页面文本与结构
Classify(分类)文档类型识别、字段抽取、规则检查文本与结构结构化字段
Output(输出)持久化、导出、通知下游结构化字段数据库/API/文件

二、核心抽象:Pipeline Stage

流水线的核心抽象是"Stage"——一个接受输入、产生输出的异步处理单元。通过统一接口,Stages可以自由组合、替换、并行化:

// 通用Stage接口
public interface IPipelineStage<TIn, TOut>
{
    string Name { get; }
    Task<TOut> ExecuteAsync(TIn input, PipelineContext ctx, CancellationToken ct);
}

// Pipeline上下文(跨阶段共享数据)
public class PipelineContext
{
    public string TraceId { get; init; } = Guid.NewGuid().ToString();
    public Dictionary<string, object> Metadata { get; } = new();
    public ILogger Logger { get; init; }
    public TelemetryScope Telemetry { get; init; }
}

// Pipeline的组合
public interface IPipeline<TIn, TOut>
{
    Task<TOut> RunAsync(TIn input, CancellationToken ct);
}

三、使用Channels实现并行流水线

.NET 6+引入的System.Threading.Channels为生产级流水线提供了高性能的背压控制和并发协调能力。下面是一个典型的三阶段并行流水线:

using System.Threading.Channels;

// 创建带背压的Channel
var preprocessChannel = Channel.CreateBounded<Document>(
    new BoundedChannelOptions(capacity: 16) {
        FullMode = BoundedChannelFullMode.Wait
    });

var extractChannel = Channel.CreateBounded<PreprocessedDoc>(16);
var outputChannel = Channel.CreateBounded<ExtractedDoc>(16);

// 启动阶段工作线程
var preprocessTasks = Enumerable.Range(0, 4).Select(_ => Task.Run(async () => {
    await foreach (var doc in preprocessChannel.Reader.ReadAllAsync())
    {
        var result = await PreprocessAsync(doc);
        await extractChannel.Writer.WriteAsync(result);
    }
})).ToArray();

var extractTasks = Enumerable.Range(0, 8).Select(_ => Task.Run(async () => {
    await foreach (var doc in extractChannel.Reader.ReadAllAsync())
    {
        var extracted = await ExtractAsync(doc);
        await outputChannel.Writer.WriteAsync(extracted);
    }
})).ToArray();

// 生产者
foreach (var path in filePaths)
{
    var doc = Document.Load(path);
    await preprocessChannel.Writer.WriteAsync(doc);
}
preprocessChannel.Writer.Complete();

// 等待链式完成
await Task.WhenAll(preprocessTasks);
extractChannel.Writer.Complete();
await Task.WhenAll(extractTasks);
outputChannel.Writer.Complete();

上述设计将流水线拆解为三级:摄取→预处理(4并发)→提取(8并发)→输出。每一级可以独立伸缩并发度,根据CPU/IO特性匹配最佳资源配比。

四、TPL Dataflow:声明式流水线

对于更复杂的流水线(带分支、合并、广播),TPL Dataflow提供了更高层的抽象:

using System.Threading.Tasks.Dataflow;

var ingest = new TransformBlock<string, Document>(
    path => Document.Load(path),
    new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 2 });

var preprocess = new TransformBlock<Document, PreprocessedDoc>(
    async doc => await PreprocessAsync(doc),
    new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 4 });

var extract = new TransformBlock<PreprocessedDoc, ExtractedDoc>(
    async doc => await ExtractAsync(doc),
    new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 8 });

var classify = new TransformBlock<ExtractedDoc, ClassifiedDoc>(
    doc => ClassifyAsync(doc),
    new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 4 });

var output = new ActionBlock<ClassifiedDoc>(
    async doc => await PersistAsync(doc),
    new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 2 });

// 连接
var linkOpts = new DataflowLinkOptions { PropagateCompletion = true };
ingest.LinkTo(preprocess, linkOpts);
preprocess.LinkTo(extract, linkOpts);
extract.LinkTo(classify, linkOpts);
classify.LinkTo(output, linkOpts);

// 投递并等待
foreach (var path in filePaths) await ingest.SendAsync(path);
ingest.Complete();
await output.Completion;

五、错误处理策略

大规模文档处理系统中,错误是常态而非例外。合理的错误处理策略应满足:

  • 单个文档失败不影响批处理整体
  • 错误原因可追溯,便于排查
  • 瞬态错误自动重试
  • 永久性错误进入死信队列
  • 错误统计可监控和告警

Result模式

在阶段间传递时使用Result<T>包装成功与失败,避免异常跨阶段传播:

public record Result<T>(bool IsSuccess, T? Value, string? ErrorCode, string? ErrorMessage)
{
    public static Result<T> Ok(T value) => new(true, value, null, null);
    public static Result<T> Fail(string code, string msg) => new(false, default, code, msg);
}

public async Task<Result<ExtractedDoc>> ExtractStageAsync(PreprocessedDoc doc)
{
    try
    {
        var extracted = await _ocr.RecognizeAsync(doc);
        return Result<ExtractedDoc>.Ok(extracted);
    }
    catch (OcrEngineException ex) when (ex.IsTransient)
    {
        return Result<ExtractedDoc>.Fail("OCR_TRANSIENT", ex.Message);
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "Extract failed for {DocId}", doc.Id);
        return Result<ExtractedDoc>.Fail("OCR_PERMANENT", ex.Message);
    }
}

六、重试与断路器

瞬态错误(网络超时、临时资源不足)应自动重试。Polly是.NET中最成熟的重试与弹性库:

using Polly;
using Polly.Retry;

var retryPolicy = new ResiliencePipelineBuilder()
    .AddRetry(new RetryStrategyOptions {
        ShouldHandle = new PredicateBuilder().Handle<OcrEngineException>(ex => ex.IsTransient),
        MaxRetryAttempts = 3,
        Delay = TimeSpan.FromSeconds(1),
        BackoffType = DelayBackoffType.Exponential,
        UseJitter = true
    })
    .AddCircuitBreaker(new CircuitBreakerStrategyOptions {
        FailureRatio = 0.5,
        MinimumThroughput = 10,
        BreakDuration = TimeSpan.FromSeconds(30)
    })
    .Build();

var result = await retryPolicy.ExecuteAsync(async ct =>
    await _ocr.RecognizeAsync(doc, ct));

七、可观测性

生产级流水线必须具备完整的可观测性:每个文档的TraceId贯穿所有阶段,每个阶段的耗时、成功率、队列长度都应可被监控。

// 使用 System.Diagnostics.ActivitySource 做分布式追踪
private static readonly ActivitySource _activity = new("DocCore.Pipeline");
private static readonly Meter _meter = new("DocCore.Pipeline");

private readonly Counter<long> _processed = _meter.CreateCounter<long>("doc_processed_total");
private readonly Histogram<double> _latency = _meter.CreateHistogram<double>("doc_stage_duration_ms");

public async Task<TOut> ExecuteAsync(TIn input, PipelineContext ctx, CancellationToken ct)
{
    using var act = _activity.StartActivity(Name);
    act?.SetTag("trace.id", ctx.TraceId);
    
    var sw = Stopwatch.StartNew();
    try
    {
        var result = await ProcessCoreAsync(input, ctx, ct);
        _processed.Add(1, new("stage", Name), new("status", "ok"));
        return result;
    }
    catch
    {
        _processed.Add(1, new("stage", Name), new("status", "error"));
        throw;
    }
    finally
    {
        _latency.Record(sw.Elapsed.TotalMilliseconds, new("stage", Name));
    }
}

八、幂等性与断点续跑

批量处理数万份文档时,不可避免会出现服务重启、节点故障。设计流水线时应考虑:

  • 幂等性:同一文档重复处理不会产生副作用(使用内容哈希作为ID)
  • 检查点:每个阶段成功后记录状态到数据库,重启时从最近检查点恢复
  • 分片:大批量任务按分片处理,已完成分片不重复
  • 死信队列:超过最大重试次数的文档进入DLQ,人工检视

九、DocCore SDK中的流水线实现

DocCore SDK内置了一个开箱即用的文档处理流水线DocumentPipeline,开发者只需配置阶段就能运行:

using DocCore.Pipeline;

// 构建流水线
var pipeline = new DocumentPipelineBuilder()
    .AddIngestStage(opts => opts.SupportedFormats = new[] { ".pdf", ".tif", ".jpg" })
    .AddPreprocessStage(opts => {
        opts.AutoDeskew = true;
        opts.Binarize = true;
        opts.TargetDpi = 300;
    })
    .AddOcrStage(opts => opts.Languages = "chi_sim+eng")
    .AddClassifyStage(opts => opts.UseRuleBook("bidding-docs.rules.json"))
    .AddOutputStage(opts => opts.JsonSink = "/data/results"))
    .WithConcurrency(preprocess: 4, ocr: 8, output: 2)
    .WithRetryPolicy(maxAttempts: 3, backoff: TimeSpan.FromSeconds(2))
    .WithDeadLetterQueue("/data/dlq")
    .Build();

// 执行
await pipeline.RunAsync(
    inputDirectory: "/data/incoming",
    progress: new Progress<PipelineProgress>(p =>
        Console.WriteLine($"处理进度 {p.Completed}/{p.Total}")),
    cancellation: cts.Token);

十、实战建议

  • 从简单开始:先用Channels实现线性流水线,再逐步引入TPL Dataflow处理复杂分支
  • CPU密集与IO密集分离:OCR是CPU密集(应限制在CPU核心数),文件读写是IO密集(可以高并发)
  • 度量优于猜测:通过指标采集找出瓶颈,不要凭直觉调参
  • 小批量验证:新流水线先用100份文档验证全链路,再推广到大批量
  • 版本化Pipeline:每次Pipeline配置变更都产生新版本,便于回溯结果差异
  • 资源回收:每个阶段的using/IDisposable必须被释放,否则大批量处理会OOM

总结

文档智能处理流水线的好坏,直接决定了系统在高负载下的吞吐量、稳定性和可维护性。.NET生态提供了从Channels到TPL Dataflow的多种并行原语,配合Polly、OpenTelemetry等组件可以构建企业级流水线。DocCore SDK封装了常见的阶段和最佳实践,开发者可以按需定制或直接使用默认配置,快速搭建文档处理系统。

想要试用DocCore SDK?

联系我们获取授权
处理流水线