1. 医疗问诊工作流系统概述
在医疗信息化快速发展的今天,如何高效、准确地处理患者问诊流程成为医疗机构面临的重要挑战。传统的手工记录和线性处理方式已经无法满足现代医疗服务的需求。本文将详细介绍一个基于LangGraph状态图的智能医疗问诊工作流系统,该系统通过模块化设计和智能体协作,实现了从患者症状输入到最终治疗建议的全流程自动化管理。
这个系统的核心价值在于将复杂的医疗问诊流程分解为多个标准化步骤,每个步骤由专门的智能体负责处理,并通过状态图精确控制流程走向。系统特别设计了状态持久化机制,使得工作流可以在任意环节中断后恢复执行,这对于需要等待化验结果的医疗场景尤为重要。同时,所有处理结果都会实时记录到数据库中,确保数据的完整性和可追溯性。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 系统架构与核心组件
2.1 整体架构设计
该系统采用分层架构设计,主要包含以下核心组件:
- 状态管理层:由MedicalWorkflowState类实现,负责维护工作流执行过程中的所有状态数据
- 工作流引擎层:MedicalWorkflow类作为核心控制器,构建并管理状态图的执行
- 智能体服务层:包含症状分析、诊断分析、治疗建议等多个专业智能体
- 数据持久层:通过Oracle数据库存储问诊记录和各环节分析结果
- 异步通信层:使用RabbitMQ处理化验结果的异步返回
这种架构设计充分考虑了医疗场景的特殊性:
- 状态持久化确保意外中断后可以恢复
- 模块化设计便于单独升级某个智能体
- 异步机制适应外部系统(如化验室)的响应延迟
2.2 MedicalWorkflowState类详解
作为系统的核心数据载体,MedicalWorkflowState类定义了工作流中的所有关键状态属性:
python复制class MedicalWorkflowState:
def __init__(self):
self.inquiry_record_id: Optional[int] = None # 问诊记录主键
self.patient_id: Optional[int] = None # 患者ID
self.doctor_id: Optional[int] = None # 医生ID
self.symptom_description: Optional[str] = None # 原始症状描述
self.symptom_analysis_result: Optional[Dict[str, Any]] = None # 症状分析结果
self.lab_tests_needed: List[str] = [] # 所需化验项目列表
self.lab_results: List[Dict[str, Any]] = [] # 化验结果集合
self.diagnosis_result: Optional[Dict[str, Any]] = None # 诊断结果
self.treatment_advice: Optional[Dict[str, Any]] = None # 治疗建议
self.current_status: InquiryStatus = InquiryStatus.INITIAL # 当前状态
self.error_message: Optional[str] = None # 错误信息
self.workflow_completed: bool = False # 完成标志
状态类的设计遵循了几个重要原则:
- 原子性:每个属性对应一个明确的业务概念
- 可序列化:提供to_dict/from_dict方法支持状态持久化
- 状态完整性:包含流程控制所需的全部信息
- 错误处理:专门的错误记录字段
提示:在实际开发中,建议为状态类添加数据校验逻辑,确保关键字段在特定状态下不为空。例如,当current_status为DIAGNOSIS_ANALYSIS时,symptom_analysis_result必须已存在。
3. 工作流引擎实现
3.1 状态图构建
MedicalWorkflow类通过StateGraph构建工作流的状态转移图:
python复制def _build_workflow_graph(self):
workflow = StateGraph(MedicalWorkflowState)
# 添加节点
workflow.add_node("create_inquiry", self._create_inquiry_record)
workflow.add_node("symptom_analysis", self._perform_symptom_analysis)
workflow.add_node("check_lab_needed", self._check_lab_tests_needed)
workflow.add_node("wait_for_lab", self._wait_for_lab_results)
workflow.add_node("diagnosis_analysis", self._perform_diagnosis_analysis)
workflow.add_node("treatment_advice", self._perform_treatment_advice)
workflow.add_node("complete_workflow", self._complete_workflow)
workflow.add_node("handle_error", self._handle_error)
# 设置边
workflow.set_entry_point("create_inquiry")
workflow.add_edge("create_inquiry", "symptom_analysis")
workflow.add_edge("symptom_analysis", "check_lab_needed")
# 条件分支
workflow.add_conditional_edges(
"check_lab_needed",
self._should_wait_for_lab,
{"wait": "wait_for_lab", "continue": "diagnosis_analysis"}
)
# 其他边
workflow.add_edge("wait_for_lab", "diagnosis_analysis")
workflow.add_edge("diagnosis_analysis", "treatment_advice")
workflow.add_edge("treatment_advice", "complete_workflow")
workflow.add_edge("complete_workflow", END)
workflow.add_edge("handle_error", END)
return workflow.compile(checkpointer=self.memory)
状态图的设计体现了医疗问诊的标准流程:
- 创建问诊记录 → 症状分析 → 检查是否需要化验
- 需要化验则进入等待状态,否则直接诊断
- 诊断分析 → 治疗建议 → 完成
3.2 关键节点实现
3.2.1 症状分析节点
python复制def _perform_symptom_analysis(self, state: MedicalWorkflowState) -> MedicalWorkflowState:
try:
request = SymptomAnalysisRequest(
patient_id=state.patient_id,
doctor_id=state.doctor_id,
symptom_description=state.symptom_description
)
result = symptom_analysis_agent.analyze_symptoms(request)
if result["success"]:
state.symptom_analysis_result = result["analysis_result"]
state.lab_tests_needed = result["lab_tests_needed"]
# 保存到数据库
analysis_result = SymptomAnalysisResult(**result["analysis_result"])
symptom_analysis_agent.save_analysis_result(
inquiry_record_id=state.inquiry_record_id,
patient_id=state.patient_id,
analysis_result=analysis_result
)
state.current_status = InquiryStatus.SYMPTOM_ANALYSIS
else:
raise Exception(result["error"])
except Exception as e:
state.error_message = str(e)
state.current_status = InquiryStatus.FAILED
return state
症状分析节点的关键点:
- 构建标准化的请求对象
- 调用专业智能体进行分析
- 处理智能体返回的结构化结果
- 实时保存分析结果到数据库
- 完善的错误处理和状态更新
3.2.2 化验等待节点
python复制def _wait_for_lab_results(self, state: MedicalWorkflowState) -> MedicalWorkflowState:
# 实际实现中会订阅RabbitMQ队列
logger.info(f"等待化验结果,问询记录ID: {state.inquiry_record_id}")
# 更新数据库状态
session = next(get_db_session())
InquiryRecordDAO.update_inquiry_status(
session=session,
inquiry_id=state.inquiry_record_id,
status=InquiryStatus.WAITING_LAB
)
session.close()
return state
这个节点的特殊之处在于:
- 工作流会在此暂停执行
- 外部系统通过RabbitMQ异步通知化验结果
- 需要专门的恢复机制继续工作流
实操技巧:在实际部署时,建议为等待状态设置超时机制(如30分钟),避免因化验系统故障导致工作流长期挂起。
4. 数据库设计与优化
4.1 核心表结构
系统使用Oracle数据库存储关键业务数据,主要表包括:
-
INQUIRY_RECORDS(问诊记录表)
- ID: 主键
- PATIENT_ID: 患者ID
- DOCTOR_ID: 医生ID
- STATUS: 当前状态
- CREATED_AT: 创建时间
- UPDATED_AT: 更新时间
-
SYMPTOM_ANALYSIS_RESULTS(症状分析结果表)
- INQUIRY_ID: 外键
- ANALYSIS_DATA: 分析结果(JSON)
- CREATED_AT: 创建时间
-
DIAGNOSIS_RESULTS(诊断结果表)
- INQUIRY_ID: 外键
- DIAGNOSIS_DATA: 诊断结果(JSON)
- CREATED_AT: 创建时间
4.2 性能优化实践
针对医疗系统的高并发需求,我们采取了以下优化措施:
-
索引优化:
- 为所有外键字段创建索引
- 为INQUIRY_RECORDS表的STATUS字段创建位图索引
- 复合索引(PATIENT_ID, CREATED_AT)支持患者历史查询
-
分区策略:
- 按时间范围分区(每月一个分区)
- 热点表采用哈希分区分散IO压力
-
连接池配置:
python复制from sqlalchemy import create_engine from sqlalchemy.pool import QueuePool engine = create_engine( "oracle+cx_oracle://user:pass@host:1521/dbname", poolclass=QueuePool, pool_size=20, max_overflow=10, pool_timeout=30 ) -
批量操作:
- 使用executemany进行批量插入
- 夜间批量处理统计分析任务
5. 异常处理与系统监控
5.1 错误处理机制
系统采用分层错误处理策略:
-
节点级错误处理:
python复制def _perform_diagnosis_analysis(self, state): try: # 业务逻辑 except AnalysisError as e: logger.error(f"诊断分析业务异常: {e}") state.error_message = f"分析异常: {e}" state.current_status = InquiryStatus.FAILED except DBError as e: logger.error(f"数据库操作失败: {e}") state.error_message = "系统错误,请稍后重试" state.current_status = InquiryStatus.RETRY except Exception as e: logger.exception("未预期的异常") state.error_message = "系统内部错误" state.current_status = InquiryStatus.FAILED return state -
工作流级错误处理:
- 通过handle_error节点统一记录错误
- 支持自动重试特定类型的错误
-
全局异常捕获:
- 使用Sentry监控系统异常
- 关键操作添加事务回滚
5.2 监控指标设计
为确保系统稳定运行,我们定义了以下监控指标:
-
性能指标:
- 各节点平均处理时间
- 工作流完成率
- 数据库查询耗时
-
业务指标:
- 每日问诊量
- 各状态工作流分布
- 化验等待时长
-
错误指标:
- 各节点错误率
- 错误类型分布
- 重试成功率
监控面板示例:
sql复制SELECT
COUNT(*) as total_inquiries,
SUM(CASE WHEN status = 'COMPLETED' THEN 1 ELSE 0 END) as completed,
AVG(EXTRACT(EPOCH FROM (updated_at - created_at))) as avg_duration
FROM inquiry_records
WHERE created_at >= TRUNC(SYSDATE)
6. 系统部署与扩展
6.1 容器化部署
系统采用Docker容器化部署,主要服务包括:
- 工作流服务:运行MedicalWorkflow核心引擎
- 智能体服务:各专业分析智能体
- 数据库服务:Oracle数据库集群
- 消息队列服务:RabbitMQ集群
docker-compose.yml关键配置:
yaml复制version: '3'
services:
workflow-engine:
image: medical-workflow:latest
environment:
- DB_URL=oracle://user:pass@db:1521/ORCL
- RABBITMQ_URL=amqp://rabbitmq
depends_on:
- db
- rabbitmq
symptom-agent:
image: symptom-analysis:2.1
environment:
- MODEL_PATH=/models/symptom-v3.h5
db:
image: oracle/database:19.3
volumes:
- oracle_data:/ORCL
6.2 水平扩展策略
针对不同组件的扩展需求:
-
无状态服务(如智能体):
- 直接增加实例数量
- 通过负载均衡分发请求
-
有状态服务(如工作流引擎):
- 采用共享存储保存状态
- 使用Redis缓存活跃状态
- 基于业务分区(如按科室划分)
-
数据库扩展:
- 读写分离
- 分库分表
- 使用Oracle RAC集群
7. 实际应用中的经验总结
在半年多的生产环境运行中,我们积累了以下宝贵经验:
-
状态序列化优化:
- 初始使用pickle序列化发现性能瓶颈
- 改用MessagePack后体积减少60%
- 最终采用Protocol Buffers实现最佳性能
-
智能体版本管理:
python复制# 智能体工厂支持多版本共存 def get_symptom_agent(version='default'): if version == 'v2': return SymptomAnalysisAgentV2() elif version == 'legacy': return LegacySymptomAgent() else: return SymptomAnalysisAgent() -
工作流版本迁移:
- 设计向后兼容的状态结构
- 开发状态迁移工具
- 双跑验证确保数据一致性
-
压力测试发现的问题:
- 数据库连接泄漏 → 引入连接池监控
- 消息堆积 → 增加消费者数量
- 状态冲突 → 优化锁策略
医疗问诊工作流系统的开发是一个持续优化的过程。随着业务需求的变化和技术的进步,我们仍在不断改进系统的各个方面。建议新用户在实施类似系统时,先从核心流程开始,逐步添加高级功能,同时建立完善的监控体系,确保系统稳定可靠地运行。
