Dask处理临床文本数据遇GC警告后挂起,求优化方案
优化方案
一、数据分区优化
- 合并过小分区:5万+分区会导致Dask调度开销激增,且频繁创建销毁对象触发GC。直接按文件大小合并成每个分区100-200MB(临床文本Parquet压缩后这个量级适配性最好):
import dask.dataframe as dd df = dd.read_parquet("path/to/year/month/*.parquet") # 按目标大小合并分区 df = df.repartition(partition_size="150MB") # 单月数据也可手动指定分区数,比如设为20-30个 # df = df.repartition(npartitions=25) - 提前过滤无用字段:只加载文本字段和必要标识,砍掉冗余数据减少内存占用:
df = dd.read_parquet("path/to/year/month/*.parquet", columns=["note_text", "label"])
二、内存与GC策略优化
- 给Dask Worker设内存上限:共享机器下避免内存溢出触发频繁GC,启动Client时配置:
from dask.distributed import Client # 64核机器按1TB总内存分配,给每个Worker留12GB(扣除系统和其他进程占用) client = Client(n_workers=64, threads_per_worker=1, memory_limit="12GB") - 开启Worker自动内存释放:让Worker在内存占用达阈值时自动清理未使用对象:
client.run(lambda dask_worker: dask_worker.memory_manager.target_fraction = 0.8) - 全局加载SpaCy模型:每个Worker只初始化一次模型,避免任务重复加载占内存:
import spacy from dask.distributed import get_worker from textdescriptives import extract_df def load_nlp_model(): worker = get_worker() if not hasattr(worker, "nlp"): # 用临床轻量模型,别加载大模型 worker.nlp = spacy.load("en_core_sci_sm") return worker.nlp def process_text(text): nlp = load_nlp_model() doc = nlp(text) return extract_df(doc).iloc[0].to_dict() df["features"] = df["note_text"].apply(process_text, meta=object)
三、预处理流程优化
- 用轻量SpaCy组件:加载模型时剔除不需要的组件(比如解析器、NER),减少内存占用:
def load_nlp_model(): worker = get_worker() if not hasattr(worker, "nlp"): nlp = spacy.load("en_core_sci_sm", exclude=["parser", "ner", "lemmatizer"]) return worker.nlp - 拆分预处理步骤:先做轻量文本清洗,再做特征提取,降低单个任务的内存负载:
def clean_text(text): text = text.strip() text = " ".join(text.split()) return text df["clean_note"] = df["note_text"].apply(clean_text, meta=str) df["features"] = df["clean_note"].apply(process_text, meta=object)
四、Dask调度与执行优化
- 切换FIFO调度器:共享机器下FIFO比自适应调度器更稳定,减少调度开销:
client = Client(scheduler="fifo") - 用
map_partitions批量处理:代替逐行apply,一次性处理整个分区,减少任务数量:def process_partition(df_part): nlp = load_nlp_model() df_part["features"] = df_part["clean_note"].apply(lambda x: extract_df(nlp(x)).iloc[0].to_dict()) return df_part df = df.map_partitions(process_partition, meta=df.dtypes.append(pd.Series([object], index=["features"])))
五、计算策略调整
- 分阶段持久化:先完成文本清洗,把结果存到Parquet,再加载做后续特征提取,避免重复计算:
df_clean = df["clean_note"].to_frame() df_clean.to_parquet("path/to/cleaned_data.parquet", overwrite=True) df_clean = dd.read_parquet("path/to/cleaned_data.parquet") - 规则过滤减少推理量:针对占比超60%的Missing类,先通过简单规则过滤(比如短文本、无临床术语的直接标记),减少模型推理样本:
def rule_based_filter(text): if len(text) < 50 or "no findings" in text.lower(): return "Missing" return None df["pred_rule"] = df["clean_note"].apply(rule_based_filter, meta=str) # 只对规则未命中的样本做模型推理 df_to_predict = df[df["pred_rule"].isna()]
内容的提问来源于stack exchange,提问作者shaun
相关产品推荐
相关产品推荐

