为什么需要流水线架构?
在文档智能处理场景中,一份PDF文件从进入系统到产生结构化结果,往往要经历十几个处理环节:解压、预处理、页面渲染、OCR识别、文本清洗、字段提取、分类、校验、输出存储等。如果将所有步骤串行硬编码在一个方法中,代码难以维护、无法并行化、出错时必须全链路重试。流水线(Pipeline)架构通过将处理分解为一系列独立、可组合的阶段(Stage),每个阶段只关心单一职责,既提升了可维护性,也释放了并发能力。
本文系统讲解如何在.NET中设计一个生产级的文档处理流水线,涵盖架构设计、阶段划分、异步模式、错误处理、重试与监控,并给出DocCore SDK中实际使用的代码样例。
一、流水线的标准阶段
一个典型的文档智能处理流水线包含以下五个核心阶段:
| 阶段 | 主要任务 | 输入 | 输出 |
|---|---|---|---|
| Ingest(摄取) | 接入文件、检测类型、基础元数据 | 文件路径/流 | 原始文档对象 |
| Preprocess(预处理) | 页面渲染、图像增强、倾斜校正 | 原始文档 | 规范化页面 |
| Extract(提取) | OCR、表格识别、布局分析 | 规范化页面 | 文本与结构 |
| Classify(分类) | 文档类型识别、字段抽取、规则检查 | 文本与结构 | 结构化字段 |
| Output(输出) | 持久化、导出、通知下游 | 结构化字段 | 数据库/API/文件 |
二、核心抽象:Pipeline Stage
流水线的核心抽象是"Stage"——一个接受输入、产生输出的异步处理单元。通过统一接口,Stages可以自由组合、替换、并行化:
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为生产级流水线提供了高性能的背压控制和并发协调能力。下面是一个典型的三阶段并行流水线:
// 创建带背压的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提供了更高层的抽象:
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 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.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贯穿所有阶段,每个阶段的耗时、成功率、队列长度都应可被监控。
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,开发者只需配置阶段就能运行:
// 构建流水线
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?
联系我们获取授权