用langchain构建数据清洗自动化流水线,将缺失值填充、去重、标准化等封装为可复用runnable组件,支持类型校验、lcel链式组合、动态配置、外部服务接入及结果验证分流。
☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 多模态理解力帮你轻松跨越从0到1的创作门槛☜☜☜

用LangChain构建数据处理自动化清洗流,不是写一堆独立函数再手动串联,而是把清洗逻辑封装成可复用、可配置、可调试的Runnable组件,让缺失值填充、重复项剔除、格式标准化等操作像流水线一样自动触发并传递上下文。
定义清洗任务的Runnable接口
第一步,明确你要清洗的数据形态:是DataFrame还是纯文本?如果是结构化数据,直接用pandas操作;如果是文档类文本(PDF/Word),先用langchain.document_loaders加载为Document对象。不要跳过这一步——【传入类型不匹配会导致后续所有Runnable.invoke失败且无明确报错】。
第二步,用RunnableLambda封装基础清洗函数。例如去除空白行和多余空格:
from langchain_core.runnables import RunnableLambdaclean_text = RunnableLambda(lambda x: x.strip().replace("\n\n", "\n").replace(" ", " "))
这个函数不能直接处理Document列表,必须先map或batch调用。若你传入的是Document对象列表,得先用RunnableParallel或自定义map逻辑做预处理。
组合多阶段清洗链(LCEL语法)
清洗不是单点动作,而是有依赖顺序的流程:去重→标准化→验证→分块。LangChain推荐用LCEL表达式链式组装,而不是嵌套函数调用。
方法一:用 | 符号串联(推荐用于线性流程)
from langchain_core.output_parsers import StrOutputParserfrom langchain_community.document_loaders import TextLoaderloader = TextLoader("raw.txt")clean_chain = ( loader.load | (lambda docs: [d.page_content for d in docs]) | clean_text | (lambda text: text.split("。")) | (lambda sentences: [s.strip() for s in sentences if s.strip()]))
注意:中间任意一步返回None或空列表,后续步骤会直接中断。比如split("。")在无句号文本中返回单元素列表["原始文本"],但若原始文本为空字符串,就会得到[""],经strip后变成[""]→过滤后为空列表→下游拿不到输入。
注入动态配置与条件分支
真实业务中,清洗规则随数据源变化:销售日志要删“测试订单”,用户反馈要保留“投诉”关键词。LangChain支持运行时传参,避免硬编码。
第一步:定义带config参数的Runnable
def conditional_clean(text: str, config: dict) -> str: if config.get("drop_test_records") and "测试订单" in text: return "" return text.replace(config.get("replace_char", " "), "")
第二步:用with_config绑定默认配置
dynamic_cleaner = RunnableLambda(conditional_clean).with_config( {"run_name": "sales_cleaner", "drop_test_records": True})
第三步:调用时覆盖config
result = dynamic_cleaner.invoke("测试订单#20260728", config={"drop_test_records": False})
这一步必须显式传入config字典,否则不会生效。LangChain不会自动合并顶层config和局部config,【未传config时,with_config设置的值不会自动注入】。
接入外部清洗服务(如视频流净化)
当本地清洗能力不足时(如去视频水印、OCR纠错),LangChain允许将清洗逻辑外置为HTTP服务或CLI工具,保持主流程干净。
第一步:封装远程调用为Runnable
import subprocessdef call_video_cleaner(video_path: str) -> str: cmd = f"python video_cleaner.py --input {video_path} --output ./clean/" subprocess.run(cmd, shell=True, check=True) return f"./clean/{os.path.basename(video_path)}"video_clean = RunnableLambda(call_video_cleaner)
第二步:与文档加载器组合
from langchain.document_loaders import UnstructuredVideoLoadervideo_loader = UnstructuredVideoLoader("raw.mp4")full_video_pipeline = video_loader.load | (lambda x: x[0].metadata["source"]) | video_clean
注意:UnstructuredVideoLoader返回的是Document对象,其.metadata["source"]才是原始路径字符串。直接传Document会触发subprocess错误。
验证清洗结果并触发告警
清洗不是终点,而是质量控制起点。LangChain支持在链末端插入校验逻辑,并根据结果分流。
步骤一:定义校验函数
def validate_cleaned(text: str) -> dict: length_ok = len(text) > 10 no_placeholder = "N/A" not in text and "暂无" not in text return {"valid": length_ok and no_placeholder, "text": text}
步骤二:用RunnableBranch做条件路由
from langchain_core.runnables import RunnableBranchvalidation_branch = RunnableBranch( (lambda x: x["valid"], lambda x: f"✅ 已通过: {x['text'][:50]}..."), (lambda x: True, lambda x: f"❌ 失败: {x['text']}"))
步骤三:完整链路
full_clean_flow = clean_text | validate_cleaned | validation_branch
这一步执行后,输出直接是字符串,无需额外解析。若需进一步处理失败样本,可在第二个分支里调用另一个Runnable(如存入error_queue)。










