Python多线程批量处理term_list写入CSV的超时与性能问题
解决方案:多线程批量处理term_list,异常/超时不中断整体流程
原代码问题分析
- 参数传递错误:
ThreadPoolExecutor.submit时误将term_list_batch传入,导致每个任务处理全量数据,而非单个term。 - 全局超时导致流程中断:
as_completed的全局超时会直接抛出TimeoutError并终止后续处理,无法继续处理未完成的正常任务。 - 线程不安全操作:
run_mappers中直接修改全局from_mapper列表、写入csv_writer,无线程锁保护,会引发数据竞争,导致CSV写入乱序或数据损坏。 - 无效的任务取消:
future.cancel()仅对未启动的任务有效,无法终止已在运行的线程。
改进方案
1. 修复线程安全问题
为csv_writer和from_mapper添加线程锁,确保多线程环境下的操作原子性。
2. 单个任务独立超时处理
放弃全局超时,改为对每个future.result()设置超时,单个任务超时/异常时仅跳过该任务,不影响其他任务执行。
3. 正确提交任务
每个任务仅传入单个term,而非全量批次。
4. 优化资源管理
合理控制线程数,按需触发GC,避免内存占用过高。
修改后的完整代码
import concurrent.futures import gc import time import threading # 全局锁,保护csv_writer和from_mapper的线程安全操作 csv_lock = threading.Lock() list_lock = threading.Lock() from_mapper = [] def run_mappers(individual_string, other_args, csv_writer): # 替换为实际耗时处理逻辑 processed_result = [individual_string, f"processed_{individual_string}"] + other_args # 线程安全地更新列表和写入CSV with list_lock: from_mapper.append(processed_result) with csv_lock: csv_writer.writerow(processed_result) return processed_result def parallelize_mappers(term_list_batch, other_args, csv_writer, max_workers=6, task_timeout=120): future_to_term = {} with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor: # 正确提交每个term的任务 future_to_term = { executor.submit(run_mappers, term, other_args, csv_writer): term for term in term_list_batch } processed_count = 0 total_terms = len(term_list_batch) for future in concurrent.futures.as_completed(future_to_term): term = future_to_term[future] processed_count += 1 try: # 单个任务设置超时 result = future.result(timeout=task_timeout) # 可选:对结果做额外处理 except concurrent.futures.TimeoutError: print(f"任务处理超时:{term}") except Exception as exc: print(f"任务执行异常 {term}: {str(exc)}") finally: # 每处理10个任务触发一次GC,控制内存 if processed_count % 10 == 0: gc.collect() time.sleep(0.5) # 短暂停顿,降低CPU负载 print(f"批次处理完成:成功处理 {processed_count} 个任务,总任务数 {total_terms}") # 测试代码 if __name__ == "__main__": import csv term_list = [ 'Dementia', 'HER2-positive Breast Cancer', 'Stroke', 'Hemiplegia', 'Type 1 Diabetes', 'Hypospadias', 'IBD', 'Eating', 'Gastric Cancer', 'Lung Cancer', 'Carcinoid', 'Lymphoma', 'Psoriasis', 'Fallopian Tube Cancer', 'Endstage Renal Disease', 'Healthy', 'HRV', 'Recurrent Small Lymphocytic Lymphoma', 'Gastric Cancer Stage III', 'Amputations', 'Asthma', 'Lymphoma', 'Neuroblastoma', 'Breast Cancer', 'Healthy', 'Asthma', 'Carcinoma, Breast', 'Fractures', 'Psoriatic Arthritis', 'ALS', 'HIV', 'Carcinoma of Unknown Primary', 'Asthma', 'Obesity', 'Anxiety', 'Myeloma', 'Obesity', 'Asthma', 'Nursing', 'Denture, Partial, Removable', 'Dental Prosthesis Retention', 'Obesity', 'Ventricular Tachycardia', 'Panic Disorder', 'Schizophrenia', 'Pain', 'Smallpox', 'Trauma', 'Proteinuria', 'Head and Neck Cancer', 'C14', 'Delirium', 'Paraplegia', 'Sarcoma', 'Favism', 'Cerebral Palsy', 'Pain', 'Signs and Symptoms, Digestive', 'Cancer', 'Obesity', 'FHD', 'Asthma', 'Bipolar Disorder', 'Healthy', 'Ayerza Syndrome', 'Obesity', 'Healthy', 'Focal Dystonia', 'Colonoscopy', 'ART', 'Interstitial Lung Disease', 'Schistosoma Mansoni', 'IBD', 'AIDS', 'COVID-19', 'Vaccines', 'Beliefs', 'SAH', 'Gastroenteritis Escherichia Coli', 'Immunisation', 'Body Weight', 'Nonalcoholic Steatohepatitis', 'Nonalcoholic Fatty Liver Disease', 'Prostate Cancer', 'Covid19', 'Sarcoma', 'Stroke', 'Liver Diseases', 'Stage IV Prostate Cancer', 'Measles', 'Caregiver Burden', 'Adherence, Treatment', 'Fracture of Distal End of Radius', 'Upper Limb Fracture', 'Smallpox', 'Sepsis', 'Gonorrhea', 'Respiratory Syncytial Virus Infections', 'HPV', 'Actinic Keratosis' ] # 分批处理(示例:每20个term为一批) batch_size = 20 other_args = ["extra_arg1", "extra_arg2"] with open("output.csv", "w", newline="", encoding="utf-8") as f: csv_writer = csv.writer(f) # 写入表头 csv_writer.writerow(["term", "processed_term", "arg1", "arg2"]) for i in range(0, len(term_list), batch_size): batch = term_list[i:i+batch_size] print(f"开始处理批次 {i//batch_size + 1},共 {len(batch)} 个任务") parallelize_mappers(batch, other_args, csv_writer) print("所有批次处理完成,结果已写入output.csv") print(f"总处理结果数:{len(from_mapper)}")
关键改进点说明
- 线程锁保护:通过
csv_lock和list_lock确保多线程下CSV写入和列表更新的原子性,避免数据混乱。 - 单个任务超时:在
future.result()中设置超时,单个任务超时仅打印日志,不中断整个批次的处理流程。 - 正确任务提交:每个任务仅处理单个term,避免重复处理全量数据。
- 内存控制:每处理10个任务触发一次GC,结合短暂停顿,平衡性能与内存占用。
- 批次化处理:将大列表拆分为小批次,进一步控制内存占用,避免一次性加载过多任务。
内容的提问来源于stack exchange,提问作者potato
相关产品推荐
相关产品推荐

