You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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)

关键修改说明

  1. 结果收集逻辑优化:让mapper返回处理后的元组,主进程遍历pool.imap结果调用reducer,避免多进程全局变量共享问题。
  2. 异常处理修复:捕获StopIteration后返回(user_name, None),在reducer中统一处理未找到的情况,避免变量未定义错误。
  3. 进度条实现:使用tqdm包装pool.imap迭代器,total参数指定总任务数,实时显示处理进度。
  4. 简化字典取值:用dict.get()替代if in判断,代码更简洁且自动设置默认值。
  5. 添加进程保护:多进程代码必须放在if __name__ == '__main__'判断下,避免子进程重复初始化代码导致的异常。

内容的提问来源于stack exchange,提问作者Patthebug

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 21:37:52