如何并行化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,无法直接并行请求所有页面。我们可以分两步实现:
- 预获取所有游标:先串行遍历一次,收集所有页面的cursor值,确保不遗漏任何页面。
- 并行处理页面请求与解析:使用线程池(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
相关产品推荐
相关产品推荐

