When One Process Becomes Too Much: Splitting a Pipeline into MCP Services
TL;DR · AI 摘要
将紧密耦合的Python流水线拆分为独立MCP服务可解决依赖冲突、故障传播和部署问题。
核心要点
- MCP服务通过进程边界隔离避免单点故障传播
- 拆分后每个服务可独立部署更新减少停机时间
- 依赖冲突问题通过独立虚拟环境解决
结构提纲
按章节快速跳转。
思维导图
用一张图看清主题之间的关系。
查看大纲文本(无障碍 / 无 JS 友好)
- 拆分流水线为MCP服务
- 旧方法缺陷
- 异常传播风险
- 依赖冲突问题
- 部署耦合度高
- MCP解决方案
- 进程边界隔离
- 独立虚拟环境
- 独立部署周期
金句 / Highlights
值得收藏与分享的关键句。
未处理异常会导致全链路崩溃,因为缺乏进程边界隔离
pip无法识别概念上独立的服务,导致依赖冲突破坏安装
MCP的核心价值在于进程边界而非通信协议
当一个进程变得过于庞大:将流水线拆分为MCP服务 | Towards Data Science
Agentic AI
当一个进程变得过于庞大:将流水线拆分为MCP服务
我们如何将紧密耦合的Python流水线拆分为可独立部署的服务
2026年9月9日
8分钟阅读
该图片使用TDS图像生成工具生成
构建多阶段流水线的常规方式是将所有组件放在同一个进程内,通过函数调用实现组件间通信。这是最简单的实现路径,在组件规模尚小时确实可行。但当每个组件都发展成具有独立依赖、独立故障模式和独立发布节奏的完整应用时,继续共享进程就会变得不再便捷,反而成为系统崩溃的诱因。我们正是在遇到这种困境后,才为每个组件引入真正的边界——通过MCP服务器实现隔离。
···
传统方案及其失效边界
在引入统一调度器类之前,我们的架构是这样的:所有服务都在调度器内部实例化,共享同一个进程和虚拟环境:
python
# 旧方案:示例代码,非真实代码,但结构是真实的
class Orchestrator:
def __init__(self):
self.extraction = ExtractionEngine() # 自带复杂的依赖树
self.risk = RiskEngine() # 使用同一库的不同版本
self.fraud = FraudEngine() # 引入与前两者冲突的依赖
def run(self, document):
fields = self.extraction.extract(document) # 此处可能抛出未处理异常
score = self.risk.score(fields) # 这行代码永远不会执行
return score仅从这7行代码中,就隐藏着三种潜在故障模式。当提取模块抛出未处理异常时,风险评分模块也会随之崩溃,因为两者之间没有隔离边界,只是简单的调用栈关系,导致整个流水线直接中断而非逐步降级。欺诈检测模块的依赖升级可能破坏提取模块的安装,因为pip无法识别这些本应是独立服务的组件,它只看到共享的单一环境。而向任意模块推送补丁都需要重新部署整个应用,因为整个系统被视为单一应用,需要重启单一进程,所有组件都处于同一个爆炸半径内。
这些问题通常在发布数月后才显现,不会出现在演示环境中。当每个引擎都积累了足够多的实际逻辑和真实依赖后,"直接导入"的便利性就会消失。故障通常表现为某个版本锁定悄无声息地强制降级其他模块,或内存密集型的提取操作耗尽风险引擎在共享进程中的资源,或欺诈检测模块一个看似无关的单行代码修改未经过提取模块的测试。它们共享同一地址空间,这种耦合程度远超最初设计时的预期。
实际的解决方案并不是 MCP 本身,而是为每个服务设置明确的进程边界。这样即使某个服务崩溃也不会影响其他服务,某个服务的依赖变更也不会波及到其他服务,每个服务都可以按照自己的时间表独立部署、扩展和重启。通过普通的 REST 端点和手动编写客户端,你同样可以实现这种边界,许多系统正是通过这种方式成功运行。MCP 带来的额外价值在于,它为任何编排器、任何代理、任何未来的消费者提供了一种统一且一致的方式,无需为每个服务和每个消费者编写定制化集成,就能发现服务的功能并调用它。
保持这种区分非常重要,因为它能防止你因惯性而非实际需求而盲目使用 MCP。如果某个服务只会有单一调用者且这个情况永远不会改变,直接使用内部 API 会更简单,协议层在这种情况下不会带来任何价值。MCP 的适用场景是特定的:当编排器需要一种统一的方式去访问越来越多的服务,并在不同条件下调用不同的服务时,这才是使用 MCP 的真实理由。
两个服务器刻意保持完全独立
编写 MCP 工具本身并不是难点;只需装饰一个函数、为参数添加类型提示,FastMCP 就会自动处理模式定义。如果你对这种机制还不熟悉,无论工具具体实现什么功能,这三行代码的结构都是一样的。真正具有挑战的部分,也是即使你了解语法也容易出错的地方,是刻意从每个服务器中排除的内容,以确保隔离性是实质性的而非表面的。
以下是提取服务的实现:
# extraction_server.py
from fastmcp import FastMCP
mcp = FastMCP("extraction-tools")
class ExtractionFailedError(Exception):
pass
@mcp.tool
def extract_fields(document: str) -> dict:
"""从原始文档中提取结构化字段"""
result = extraction_engine.run(document)
if result.confidence < 0.6:
raise ExtractionFailedError(
f"提取置信度 {result.confidence:.2f} 低于阈值"
)
return result.fields
if __name__ == "__main__":
mcp.run(
transport="http",
host="0.0.0.0",
port=8931,
)以及风险评分服务:
# risk_server.py
from fastmcp import FastMCP
mcp = FastMCP("risk-tools")
@mcp.tool
def score_risk(fields: dict) -> dict:
"""对提取的字段进行风险评分"""
return risk_engine.score(fields)
if __name__ == "__main__":
mcp.run(
transport="http",
host="0.0.0.0",
port=8932,
)两个文件之间没有任何相互引用,也没有安装彼此的依赖项。它们作为两个独立的进程运行在不同的端口上,这不是偶然的部署细节,而是这个架构存在的根本原因。如果下个季度风险引擎需要升级库版本,这只是一个进程的变更,可以独立测试和部署,而 extraction_server.py 完全不会感知到这个变化。可流式传输的 HTTP 协议是这种部署模型最合适的传输方式;而 stdio 会将服务器的生命周期绑定到启动它的进程,这虽然适合本地客户端,但不适合独立部署的服务。
什么是容易被忽视的是,边界只有在你实际编写工具时尊重它才会生效。一旦两个服务器都存在,人们很容易让 extraction_server.py 导入一个共享的 utils.py,而这个 utils.py 同时也被 risk_server.py 导入。虽然这很便利,但会悄然重新创建你原本构建两个服务器就是为了避免的耦合。如果两个服务都需要相同的验证逻辑,要么需要在每个服务中刻意重复实现该逻辑,要么让该逻辑成为拥有自己工具的第三个独立服务。这种隔离性是一种需要刻意保持的纪律,而不是使用 MCP 就能自动获得的默认特性。
仅使用 SDK 调用工具
在任何编排框架介入之前,了解客户端实际执行的操作是有价值的,因为所有构建在该基础之上的功能本质上都在调用相同的两个方法:
from mcp import ClientSession
from mcp.client.streamable_http import streamablehttp_client
async def call_extraction(raw_text: str):
async with streamablehttp_client(
"http://localhost:8931/mcp"
) as (read, write, _):
async with ClientSession(read, write) as session:
await session.initialize()
tools = await session.list_tools()
print([tool.name for tool in tools.tools])
result = await session.call_tool(
"extract_fields",
{"document": raw_text},
)
return result发现可用功能,通过名称和参数调用对应功能。其他如模式定义、传输方式、多服务器路由等机制的存在,都是为了在大规模场景下让这种模式有效运行,而不是用更复杂的方案取代它。
将服务连接到编排器
编排器的工作不是选择一个工具并调用它,而是计算服务之间的执行路径,其中每一步的输出都会成为下一步的输入。langchain-mcp-adapters 可以同时连接到多个服务器,而 LangGraph 的 StateGraph 是明确表达这种路径的合理方式:
from typing import TypedDict
from langchain_mcp_adapters.client import MultiServerMCPClient
from langgraph.graph import END, START, StateGraph
client = MultiServerMCPClient(
{
"extraction": {
"url": "http://localhost:8931/mcp",
"transport": "streamable_http"
},
"risk": {
"url": "http://localhost:8932/mcp",
"transport": "streamable_http"
}
}
)
class PipelineState(TypedDict):
document: str
fields: dict
risk_score: dict
async def extraction_node(state: PipelineState):
tools = await client.get_tools(server_name="extraction")
tool = next(t for t in tools if t.name == "extract_fields")
fields = await tool.ainvoke({"document": state["document"]})
return {"fields": fields}
async def risk_node(state: PipelineState):
tools = await client.get_tools(server_name="risk")
tool = next(t for t in tools if t.name == "score_risk")
score = await tool.ainvoke({"fields": state["fields"]})
return {"risk_score": score}
graph = StateGraph(PipelineState)
graph.add_node("extract", extraction_node)
graph.add_node("score_risk", risk_node)
graph.add_edge(START, "extract")
graph.add_edge("extract", "score_risk")
graph.add_edge("score_risk", END)
app = graph.compile()这是一个固定路径:提取操作始终会传递给风险评分,不需要任何决策。对于包含两个服务的流水线来说,这种诚实的实现方式是值得构建的,在尝试更复杂的方案之前,先明确实现这个基础版本是有价值的。当第三个或第四个服务加入,且路径确实需要根据文档内容进行动态选择时,add_edge 就不再足够,需要使用 add_conditional_edges,此时编排器会根据前一个服务的返回结果决定调用哪个服务。
这尚未解决的问题
这使你实现了隔离,并通过一个决定执行顺序的协调器,提供了一种清晰的服务间路由方式。但它并未说明当执行过程中顺序出现错误时会发生什么,例如路径中第三个步骤的服务发现某个字段缺失,而该字段只能由更早的服务提供,此时协调器必须回退而非继续执行。从提取到风险评分的固定边无法表达这种情形。它也未说明一旦服务故障变成网络调用而非三个堆栈层级上的Python异常时,故障的实际表现形式是什么;或者当协调器需要在多个服务间进行选择,而单个决策无法清晰容纳所有选项时会发生什么。这些都是独立的现实问题,它们保持独立而非被强行压缩到当前问题的结尾。
我最初提出的问题在使用MCP后已经解决。自那以后,提取和风险评分再未共享过堆栈跟踪,一个模块的依赖升级从未影响到另一个模块,向任一服务推送修复补丁也不再需要重新部署另一个服务。在解决上述更复杂问题中的任何一个之前,这种状态就已经成立,这说明原始痛点中有多少是真正关于隔离的问题,又有多少其实与协议无关。
这是以正确方式重建系统后诞生的第一个成果,但不会是最后一个。我刚刚列出的未解决问题仍有足够多的余地:执行中途出现错误的路径、通过网络表现与堆栈跟踪不同的故障、协调器需要在多个服务间选择而单个决策无法清晰容纳这些情况。单篇文章永远无法解决所有问题,我宁愿保持开放,也不愿强行给出一个比系统当前状态更整洁的结局。