Python 医学数据处理:结构化病历的标准化清洗 Pipeline
Python 医学数据处理:结构化病历的标准化清洗 Pipeline
一、医生写的病历,机器读起来像天书
电子病历(EMR)的数据质量是医疗 AI 落地最大的拦路虎。一份看似完整的病历,经过程序解析后往往会暴露各种问题:主诉字段填了 500 字的流水账、诊断名称同时存在 ICD-10 和 ICD-11 两种编码、"高血压"有时写作"高血压病"有时缩写为"HBP"。如果不做标准化清洗,下游模型训练出来就是垃圾进垃圾出。
更糟糕的是,病历数据中的错误类型五花八门:数值型字段填了文字("体温:正常"而不是"36.5")、日期格式不统一("2024/3/15"和"2024-03-15"混用)、必填字段留空但填了一个"无"字。传统 ETL 工具面对这种半结构化数据基本失灵,需要一个专门针对医学文本的清洗 Pipeline。
二、清洗 Pipeline 架构:化零为整,逐层提纯
标准化的病历清洗不是一步到位,而是分阶段提纯的过程:
每个阶段产出带标注的中间数据,这样出问题时可以快速定位是哪个环节的锅。关键设计:前一个阶段的输出不会覆盖原始数据,而是产生新列(如diagnosis_raw→diagnosis_cleaned→diagnosis_normalized),方便追溯。
三、Python 实现:可配置的清洗 Pipeline
import re import pandas as pd from datetime import datetime from typing import Dict, List, Optional, Tuple from dataclasses import dataclass, field import logging logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) @dataclass class CleaningConfig: """清洗配置""" date_formats: List[str] = field(default_factory=lambda: [ "%Y-%m-%d", "%Y/%m/%d", "%Y年%m月%d日" ]) # 症状同义词映射 symptom_synonyms: Dict[str, str] = field(default_factory=lambda: { "HBP": "高血压", "高血压病": "高血压", "血压偏高": "高血压", "DM": "糖尿病", "糖尿病Ⅱ型": "2型糖尿病", }) # 数值字段的正则 temperature_pattern: re.Pattern = re.compile( r'(\d{2}\.?\d?)\s*[°℃度]?' ) class EMRCleaner: """电子病历清洗器""" def __init__(self, config: Optional[CleaningConfig] = None): self.config = config or CleaningConfig() self.cleaning_log: List[Dict] = [] def load_raw_data(self, filepath: str) -> pd.DataFrame: """加载原始病历数据""" try: df = pd.read_csv(filepath, encoding='utf-8') except UnicodeDecodeError: df = pd.read_csv(filepath, encoding='gbk') logger.info(f"加载 {len(df)} 条病历记录") return df def standardize_dates(self, df: pd.DataFrame, date_cols: List[str]) -> pd.DataFrame: """阶段1:日期格式标准化""" for col in date_cols: if col not in df.columns: continue df[f'{col}_std'] = df[col].apply(self._parse_date) return df def _parse_date(self, value) -> Optional[str]: """尝试多种格式解析日期""" if pd.isna(value) or str(value).strip() in ('', '无', '不详'): return None value_str = str(value).strip() for fmt in self.config.date_formats: try: dt = datetime.strptime(value_str, fmt) return dt.strftime("%Y-%m-%d") except ValueError: continue # 尝试只提取数字 digits = re.findall(r'\d+', value_str) if len(digits) >= 3: return f"{digits[0]}-{digits[1].zfill(2)}-{digits[2].zfill(2)}" logger.warning(f"无法解析日期: {value_str}") return None def clean_temperature(self, df: pd.DataFrame, col: str = 'temperature') -> pd.DataFrame: """清洗体温字段:文字描述转为数值""" def extract_temp(val): if pd.isna(val): return None val_str = str(val).strip() # 处理"正常""无发热"等描述 if val_str in ('正常', '无发热', '体温正常'): return 36.5 match = self.config.temperature_pattern.search(val_str) if match: temp = float(match.group(1)) # 异常值检测 if 34 <= temp <= 43: return temp else: logger.warning(f"体温异常值: {temp},原文: {val_str}") return None return None df[f'{col}_numeric'] = df[col].apply(extract_temp) return df def normalize_symptoms(self, df: pd.DataFrame, col: str = 'chief_complaint') -> pd.DataFrame: """阶段3:症状术语归一化""" synonyms = self.config.symptom_synonyms def normalize(text: str) -> str: if pd.isna(text): return text result = str(text) for old, new in synonyms.items(): result = result.replace(old, new) return result df[f'{col}_normalized'] = df[col].apply(normalize) return df def quality_score(self, df: pd.DataFrame) -> pd.DataFrame: """阶段4:数据质量评分(0-100)""" scores = pd.Series(100.0, index=df.index) # 必填字段缺失扣分(每缺失1个字段扣15分) required_fields = ['patient_id', 'visit_date_std', 'diagnosis'] for field in required_fields: if field in df.columns: scores -= df[field].isna().astype(float) * 15 # 异常体温扣分 if 'temperature_numeric' in df.columns: temp = df['temperature_numeric'] abnormal = (temp < 35) | (temp > 42) scores -= abnormal.astype(float) * 20 # 日期未来时间扣分 if 'visit_date_std' in df.columns: future = pd.to_datetime( df['visit_date_std'], errors='coerce' ) > datetime.now() scores -= future.astype(float) * 25 scores = scores.clip(lower=0, upper=100) df['quality_score'] = scores return df def pipeline(self, filepath: str) -> Tuple[pd.DataFrame, Dict]: """执行完整清洗 Pipeline""" df = self.load_raw_data(filepath) date_cols = [c for c in df.columns if any(kw in c.lower() for kw in ['date', 'time', '日期', '时间'])] df = self.standardize_dates(df, date_cols) df = self.clean_temperature(df) df = self.normalize_symptoms(df) df = self.quality_score(df) # 生成清洗报告 report = { 'total_records': len(df), 'high_quality': int((df['quality_score'] >= 80).sum()), 'need_review': int(((df['quality_score'] >= 60) & (df['quality_score'] < 80)).sum()), 'excluded': int((df['quality_score'] < 60).sum()), 'date_parse_failures': df[[c for c in df.columns if c.endswith('_std')]] .isna().sum().to_dict(), } logger.info(f"清洗完成: {report}") return df, report四、边界分析与 Trade-offs
清洗粒度:自动 vs 人工:全自动清洗的准确率大约在 85%-90%,剩下的 10% 需要人工核实。经验法则是:质量分 >= 80 的记录直接入库,60-79 分的提交人工审核队列,< 60 分的标记后暂存(等数据源修复后重新跑)。不要在代码里硬编码太多业务规则——规则变化比代码变化快得多。
同义词映射的维护:症状和诊断的同义词是最容易过时的。建把映射表放在配置中心(如 Consul 或本地 YAML 文件),支持热更新,不要让运维为了改一个"高血压"的同义词而重新部署服务。映射表应该由医学知识团队维护,而不是开发团队。
大文件的内存策略:pandas 默认把整个 DataFrame 加载到内存,处理 10GB 的病历文件会直接 OOM。解决方案是分块读取(chunksize参数)或使用 Dask/Polars 这样的增量计算框架。分块时需要注意,某些清洗操作(如日期范围校验)需要全局上下文,这种情况可以先采样估算全局参数,再分块处理。
五、总结
病历数据清洗的核心思路是"先标准化格式,再归一化语义,最后量化质量"。分阶段处理的优势在于问题可追溯——数据质量报告会精确告诉你哪一列、哪种错误类型最严重。DataFrame 的apply操作虽然方便,但在百万级数据上会变慢,生产环境建议迁移到 Polars 或 PySpark。最重要的一点:清洗不是一次性的,数据源的格式会随时间变化(医院升级 HIS 系统等),所以 Pipeline 要设计为可重跑的,支持增量清洗和全量重洗两种模式。