"""任务:实现可配置的文档处理管道场景:灵活的文档预处理流程"""from typing import List, Callablefrom dataclasses import dataclassfrom abc import ABC, abstractmethodimport re@dataclassclass Document:"""文档对象"""id: strcontent: strmetadata: dictchunks: List[str] = Noneembeddings: List[List[float]] = Noneclass ProcessingStep(ABC):"""处理步骤基类"""@abstractmethoddef process(self, doc: Document) -> Document:passclass TextCleaningStep(ProcessingStep):"""文本清洗步骤"""def process(self, doc: Document) -> Document:text = doc.contenttext = re.sub(r'\s+', ' ',text)text = re.sub(r'[^\w\s\u4e00-\u9fff.,!?,。!?]', '', text)"""^ → 非(不要)\w → 字母数字\s → 空格\u4e00-\u9fff → 所有中文字符.,!?,。!? → 保留正常标点只保留:中文 + 英文 + 数字 + 空格 + 标准标点"""text = text.strip()doc.content = textreturn docclass ChunkingStep(ProcessingStep):"""分块步骤"""def __init__(self, chunk_size: int = 500, overlap: int = 50):self.chunk_size = chunk_sizeself.overlap = overlapdef process(self, doc: Document) -> Document:text = doc.contentchunks = []start = 0while start < len(text):end = min(start + self.chunk_size, len(text))chunks.append(text[start: end])start += self.chunk_size - self.overlapdoc.chunks = chunksreturn docclass EmbeddingStep(ProcessingStep):"""Embedding步骤"""# def __init__(self, embedding_service):# self.embedding_service = embedding_servicedef process(self, doc: Document) -> Document:if doc.chunks:doc.embeddings = [[0.1] * 10 for _ in doc.chunks]return docclass DocumentPipeline:"""文档处理管道"""def __init__(self):self.steps: List[ProcessingStep] = []def add_step(self, step: ProcessingStep):"""添加处理步骤"""self.steps.append(step)return self # 支持链式调用def process(self, doc: Document) -> Document:"""执行管道"""for step in self.steps:doc = step.process(doc)return docdef process_batch(self, docs: List[Document]) -> List[Document]:"""批量处理"""return [self.process(doc) for doc in docs]# 使用示例pipeline = (DocumentPipeline().add_step(TextCleaningStep()).add_step(ChunkingStep(chunk_size=500)).add_step(EmbeddingStep()))doc = Document(id="1", content="...", metadata={})processed = pipeline.process(doc)print(" 处理完成")print("分块数:", len(processed.chunks))print("生成向量:", processed.embeddings is not None)
一、场景
场景定位:RAG 系统通用文档预处理管道,应用于:
本地 / 上传文档统一预处理流水线:文本清洗 → 内容分块 → 生成向量
可自由增删、调换处理步骤,不用改核心流程
支持单文档处理 + 批量文档处理
适配后续接入真实向量库、真实 Embedding 模型,是可插拔的预处理骨架
完整执行流程:
1、定义数据载体,用 @dataclass 定义 Document 文档对象,统一承载:文档 ID、原文、元数据、分块列表、向量列表;chunks/embeddings 设默认值 None,实例化不用手动传,后续流程自动赋值。
2、统一抽象处理规范,基类 ProcessingStep(ABC) 定统一接口:所有处理步骤必须实现 process(doc) 方法,入参、返回都是 Document 对象。
3、实现具体处理步骤:
TextCleaningStep:正则清洗脏字符、保留中文 / 字母 / 数字 / 标准标点;
ChunkingStep:固定长度分块 + 重叠区;
EmbeddingStep:基于分块生成向量(当前是模拟向量,后续可替换真实模型)。
4、组装管道,DocumentPipeline 作为管道容器:add_step() 链式添加处理步骤,内部按添加顺序保存所有步骤。
5、执行流水线,调用 pipeline.process(doc):循环遍历每一个 step → 把上一步处理完的 doc 传给下一步 → 全程同一个 Document 对象流转、逐级赋值。
6、批量兜底,process_batch 循环调用单文档处理,支持批量文档一键走流水线。
二、代码中的关键洞察
1、管道模式核心思想,把每一个预处理动作拆成独立 Step,可插拔、可增删、可调顺序,符合开闭原则,新增处理步骤不用改管道核心代码。
2、统一接口归一化,所有清洗 / 分块 / 向量化步骤,统一继承 ProcessingStep、统一 process(doc) 入参出参,流水线不用关心具体步骤实现,只负责调度。
3、Document 作为统一数据载体,全程只靠一个 Document 对象流转,处理过程中原地给 chunks、embeddings 赋值,不用来回传多个变量,结构干净统一。
4、类型提示 ≠ 对象自动创建,doc: Document 只是语法标注,仅做可读性和编辑器提示,不会自动生成对象,必须手动实例化传入。
四、关键代码/可复用片段
1. 通用 Document 数据结构
from dataclasses import dataclassfrom typing import List@dataclassclass Document:id: strcontent: strmetadata: dictchunks: List[str] = Noneembeddings: List[List[float]] = None
2. 处理步骤统一抽象基类
from abc import ABC, abstractmethodclass ProcessingStep(ABC):@abstractmethoddef process(self, doc: Document) -> Document:pass
3. 通用文本清洗正则
import re# 保留中文、字母、数字、空格、常规标点,过滤所有特殊符号/乱码text = re.sub(r'\s+', ' ', text)text = re.sub(r'[^\w\s\u4e00-\u9fff.,!?,。!?]', '', text)text = text.strip()
4. 通用文档管道调度器(链式调用 + 批量处理)
class DocumentPipeline:def __init__(self):self.steps = []def add_step(self, step:ProcessingStep):self.steps.append(step)return self # 支持链式调用def process(self, doc:Document) -> Document:for step in self.steps:doc = step.process(doc)return docdef process_batch(self, docs:List[Document]) -> List[Document]:return [self.process(doc) for doc in docs]
5. 固定大小分块器
class ChunkingStep(ProcessingStep):def __init__(self, chunk_size: int = 500, overlap: int = 50):self.chunk_size = chunk_sizeself.overlap = overlapdef process(self, doc: Document) -> Document:text = doc.contentchunks = []start = 0while start < len(text):end = min(start + self.chunk_size, len(text))chunks.append(text[start: end])start += self.chunk_size - self.overlapdoc.chunks = chunksreturn doc
夜雨聆风