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

如何并行化OpenAlex API请求高效获取2020年论文数据?

OpenAlex API 并行化获取2020年论文实现方案

我正在使用OpenAlex API获取2020年全部论文,当前采用游标分页的串行循环方式实现,但运行速度极慢。刚接触循环并行化,需要不破坏结果完整性的并行化实现方案。

原串行代码:

import requests
from collections import defaultdict
import pandas as pd

# url with a placeholder for cursor
example_url_with_cursor ='https://api.openalex.org/works?filter=publication_year:2020&cursor={}'

dfs=defaultdict(dict)
paper_id=[]
title_lst=[]
year_lst=[]
lev_lst=[]

page=1
cursor = '*'

# loop through pages
while cursor:
    
    # set cursor value and request page from OpenAlex
    url = example_url_with_cursor.format(cursor)
    print("\n" + url)
    page_with_results = requests.get(url).json()
    
    # loop through partial list of results
    results = page_with_results['results']
    for i,work in enumerate(results):
        openalex_id = work['id'].replace("https://openalex.org/", "")
        if work['display_name'] is not None and len(work['display_name'])>0:
            openalex_title = work['display_name']
        else:
            openalex_title='No title'
        openalex_year = work['publication_year']
        if work['concepts'] is not None and  len(work['concepts'])>0:
            openalex_lev = work['concepts'][0]['display_name']
        else:
            openalex_lev = 'None'
        paper_id.append(openalex_id)
        title_lst.append(openalex_title)
        year_lst.append(openalex_year)
        lev_lst.append(openalex_lev)
        #Constructing a pandas db:
        df=pd.DataFrame(paper_id, columns=['paper_id'])
        df['title'] = title_lst
        df['level'] = lev_lst
        df['pub_year'] = year_lst

    # update cursor to meta.next_cursor
    cursor = page_with_results['meta']['next_cursor']

并行化实现思路

由于OpenAlex的游标分页依赖前一页返回的next_cursor,无法直接并行请求所有页面。我们可以分两步实现:

  1. 预获取所有游标:先串行遍历一次,收集所有页面的cursor值,确保不遗漏任何页面。
  2. 并行处理页面请求与解析:使用线程池(IO密集型任务优先选线程池)并行处理每个cursor对应的页面,解析数据后统一合并结果。

完整并行化代码

import requests
import pandas as pd
from concurrent.futures import ThreadPoolExecutor, as_completed
from collections import defaultdict

# 配置参数
BASE_URL = 'https://api.openalex.org/works?filter=publication_year:2020&cursor={}'
MAX_WORKERS = 8  # 控制并发数,避免触发API限流
TIMEOUT = 10

def get_all_cursors():
    """预获取所有页面的cursor"""
    cursors = []
    current_cursor = '*'
    while current_cursor:
        url = BASE_URL.format(current_cursor)
        try:
            response = requests.get(url, timeout=TIMEOUT)
            response.raise_for_status()
            data = response.json()
            cursors.append(current_cursor)
            current_cursor = data['meta']['next_cursor']
            print(f"已获取游标: {current_cursor}")
        except Exception as e:
            print(f"获取游标失败: {str(e)}")
            break
    return cursors

def process_page(cursor):
    """处理单个页面的请求与数据解析"""
    url = BASE_URL.format(cursor)
    page_data = []
    try:
        response = requests.get(url, timeout=TIMEOUT)
        response.raise_for_status()
        data = response.json()
        results = data.get('results', [])
        for work in results:
            # 解析单篇论文数据
            openalex_id = work['id'].replace("https://openalex.org/", "")
            openalex_title = work['display_name'] if (work.get('display_name') and len(work['display_name'])>0) else 'No title'
            openalex_year = work['publication_year']
            openalex_lev = work['concepts'][0]['display_name'] if (work.get('concepts') and len(work['concepts'])>0) else 'None'
            
            page_data.append({
                'paper_id': openalex_id,
                'title': openalex_title,
                'level': openalex_lev,
                'pub_year': openalex_year
            })
        print(f"完成页面处理: {cursor}")
    except Exception as e:
        print(f"处理页面失败 {cursor}: {str(e)}")
    return page_data

if __name__ == '__main__':
    # 第一步:预获取所有游标
    print("开始获取所有游标...")
    all_cursors = get_all_cursors()
    print(f"共获取到 {len(all_cursors)} 个页面游标")
    
    # 第二步:并行处理所有页面
    print("开始并行处理页面...")
    all_data = []
    with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor:
        # 提交所有任务
        futures = {executor.submit(process_page, cursor): cursor for cursor in all_cursors}
        # 处理完成的任务
        for future in as_completed(futures):
            page_data = future.result()
            if page_data:
                all_data.extend(page_data)
    
    # 第三步:合并为DataFrame
    print("合并数据为DataFrame...")
    final_df = pd.DataFrame(all_data)
    print(f"最终共获取 {len(final_df)} 篇论文数据")
    
    # 可选:保存到文件
    final_df.to_csv('2020_openalex_papers.csv', index=False)
    print("数据已保存到 2020_openalex_papers.csv")

关键注意事项

  • 并发数控制:MAX_WORKERS不要设置过高(建议5-10),OpenAlex API有请求频率限制,过高并发可能导致请求被拒绝或限流。
  • 异常处理:代码中加入了try-except捕获请求和解析异常,避免单个页面失败导致整个程序终止。
  • 数据安全:每个线程独立处理页面数据并返回结果列表,最后统一合并,避免多线程操作全局列表的线程安全问题。
  • 效率优化:原串行代码中每次循环都重新生成DataFrame,并行版改为最后一次性合并,大幅提升效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 12:10:22