基于multiprocessing的MapReduce获取Scholarly数据卡住问题排查
问题分析与解决方案
代码中的错误点
pool.map返回值处理错误:pool.map返回的是包含所有mapper输出的列表,直接用author, user_name = pool.map(...)会触发解包失败,因为列表长度等于用户数量而非2。- 多进程全局变量共享问题:
author_pub_details是全局列表,子进程的修改不会同步到主进程(多进程有独立内存空间),导致报错信息和结果无法正确收集。 - 异常处理逻辑漏洞:捕获
StopIteration后first_author_result未定义,后续调用RetrieveAuthorDetails(first_author_result)会触发NameError,导致进程卡住或报错。 - Reducer调用方式错误:原代码仅调用一次reducer,但
pool.map返回多个结果,需遍历每个结果调用reducer。
修正后的代码(含进度条)
from multiprocessing import Pool from scholarly import scholarly from tqdm import tqdm def GetFirstAuthor(author_name): search_query = scholarly.search_author(author_name) first_author_result = next(search_query) return first_author_result def RetrieveAuthorDetails(first_author_result): author = scholarly.fill(first_author_result) return author def RetrievePublicationDetails(author): author_pub = [] for pub in author['publications']: title = pub['bib'].get('title', '') pub_year = pub['bib'].get('pub_year', 0) author_names = pub['bib'].get('author', '') journal = pub['bib'].get('journal', '') citation = pub['bib'].get('citation', '') author_pub.append([title, pub_year, author_names, journal, citation]) return author_pub def mapper(user_name): try: first_author_result = GetFirstAuthor(user_name) author = RetrieveAuthorDetails(first_author_result) return (user_name, author) except StopIteration: return (user_name, None) # 用None标记未找到记录 def reducer(result): user_name, author = result if author is None: return [user_name, 'Record not found for the author'] else: author_pub = RetrievePublicationDetails(author) return [user_name, author_pub] if __name__ == '__main__': author_pub_details = [] user_names = ['Marc Mertens','Katharina Breininger','Marie Caroline Guzian','James M. Shwayder','M. Narasimha','Brian W. Miller','Peter C. Y. Chen','Maxim N. Peshkov'] with Pool() as pool: # 用imap结合tqdm实现进度条 for result in tqdm(pool.imap(mapper, user_names), total=len(user_names)): processed = reducer(result) author_pub_details.append(processed) # 验证结果 for item in author_pub_details: print(item)
关键修改说明
- 结果收集逻辑优化:让mapper返回处理后的元组,主进程遍历
pool.imap结果调用reducer,避免多进程全局变量共享问题。 - 异常处理修复:捕获
StopIteration后返回(user_name, None),在reducer中统一处理未找到的情况,避免变量未定义错误。 - 进度条实现:使用
tqdm包装pool.imap迭代器,total参数指定总任务数,实时显示处理进度。 - 简化字典取值:用
dict.get()替代if in判断,代码更简洁且自动设置默认值。 - 添加进程保护:多进程代码必须放在
if __name__ == '__main__'判断下,避免子进程重复初始化代码导致的异常。
内容的提问来源于stack exchange,提问作者Patthebug
相关产品推荐
相关产品推荐

