百万级嵌套JSON高效解析:提取去重无邮箱机构信息并行方案咨询
百万级嵌套JSON机构信息并行解析方案(去重+移除邮箱)
需求明确
需要处理约100万条嵌套JSON记录,提取其中的Affiliation字段内容:
- 移除字段中的电子邮件地址(格式为
Electronic address: xxx@xxx.xx) - 最终输出全局唯一的机构文本条目
- 要求用GNU Parallel或Python多进程实现并行处理,提升效率
输入示例(单条记录)
{ "_index": "group", "_type": "_doc", "_id": "9890798789", "_score": 1, "_source": { "Bibtex": { "Article": { "AuthorList": [ { "Affiliation": { "Affiliation": "Some departmentA, some university, city, country. Electronic address: gh@example.com." } }, { "Affiliation": { "Affiliation": "Some departmentB, some university, city, country. Electronic address: jh@example.com." } }, { "Affiliation": { "Affiliation": "Some departmentA, some university, city, country; Institute, Sydney, country. Electronic address: yu@example.com." } }, { "Affiliation": { "Affiliation": "Some department some university, city, country. Electronic address: nj@example.com." } }, { "Affiliation": { "Affiliation": "department, university, Sydney, country. Electronic address: bg@a.b.au." } }, { "Affiliation": { "Affiliation": "Some departmentA, some university, city, country. Electronic address: we@example.com." } } ] } } } }
输出示例(单条记录对应结果)
Some departmentA, some university, city, country; Institute, Sydney, country. Some departmentB, some university, city, country. Some department some university, city, country. department, university, Sydney, country.
方案一:GNU Parallel + jq 实现
适合熟悉命令行的场景,无需编写复杂脚本,处理速度快。
前提条件
已安装jq(JSON处理工具)和GNU Parallel:
- Debian/Ubuntu:
sudo apt install jq parallel - macOS:
brew install jq parallel
处理步骤
1. 输入为JSON Lines格式(每行一条记录)
直接用Parallel批量处理,同时完成提取、清洗、单条去重,最后全局去重:
# 用8个进程并行处理,可根据CPU核心数调整-j参数 parallel -j 8 'jq -r '"'"'._source.Bibtex.Article.AuthorList[].Affiliation.Affiliation | sub("\\. Electronic address:.*"; ".") | unique'"'"' {}' ::: input_*.json | sort -u > final_unique_affiliations.txt
2. 输入为单个超大JSON文件(非JSON Lines)
先拆分文件为小批量,避免单次处理负载过高:
# 每10000条记录拆分为一个文件,前缀split_ split -l 10000 big_input.json split_ # 并行处理拆分后的文件,最后全局去重 parallel -j 8 'jq -r '"'"'._source.Bibtex.Article.AuthorList[].Affiliation.Affiliation | sub("\\. Electronic address:.*"; ".") | unique'"'"' {}' ::: split_* | sort -u > final_unique_affiliations.txt
命令解释
jq部分:提取所有嵌套的Affiliation字段,用sub正则移除邮箱后缀,unique实现单条记录内去重parallel -j 8:启动8个并行进程sort -u:对所有结果全局去重
方案二:Python多进程实现
适合需要更灵活逻辑(比如错误处理、自定义清洗规则)的场景,支持流式处理避免内存溢出。
基础版本(JSON Lines输入)
import re import json from multiprocessing import Pool import sys def process_record(record_str): """处理单条JSON记录,返回清洗后的机构列表(单条内去重)""" try: record = json.loads(record_str) affils = set() # 逐层提取Affiliation字段 author_list = record.get("_source", {}).get("Bibtex", {}).get("Article", {}).get("AuthorList", []) for author in author_list: affil_text = author.get("Affiliation", {}).get("Affiliation", "") if not affil_text: continue # 移除邮箱部分 cleaned_affil = re.sub(r'\. Electronic address:.*', '.', affil_text).strip() if cleaned_affil: affils.add(cleaned_affil) return list(affils) except Exception: # 跳过解析失败的记录,可根据需求添加日志 return [] def main(input_file, num_workers=8): # 读取所有记录(JSON Lines格式) with open(input_file, 'r', encoding='utf-8') as f: records = [line.strip() for line in f if line.strip()] # 进程池并行处理 with Pool(num_workers) as pool: results = pool.map(process_record, records) # 全局去重 global_unique_affils = set() for res in results: global_unique_affils.update(res) # 输出结果到文件 with open("final_affiliations.txt", 'w', encoding='utf-8') as f: for affil in sorted(global_unique_affils): f.write(f"{affil}\n") if __name__ == "__main__": if len(sys.argv) != 2: print("使用方式: python process_affiliations.py <输入JSON Lines文件路径>") sys.exit(1) main(sys.argv[1])
流式处理版本(适合超大JSON数组)
如果输入是单个超大JSON数组(而非每行一条),用ijson流式解析,避免内存占用过高:
import re import ijson from multiprocessing import Pool import sys def process_batch(batch): """处理一批记录,返回去重后的机构列表""" affils = set() for record in batch: author_list = record.get("_source", {}).get("Bibtex", {}).get("Article", {}).get("AuthorList", []) for author in author_list: affil_text = author.get("Affiliation", {}).get("Affiliation", "") if not affil_text: continue cleaned_affil = re.sub(r'\. Electronic address:.*', '.', affil_text).strip() if cleaned_affil: affils.add(cleaned_affil) return list(affils) def main(input_file, num_workers=8, batch_size=1000): # 流式读取JSON数组(假设顶级结构是数组,每个元素为一条记录) batches = [] current_batch = [] with open(input_file, 'r', encoding='utf-8') as f: # ijson.items的第二个参数'item'对应数组中的每个元素 for record in ijson.items(f, 'item'): current_batch.append(record) if len(current_batch) >= batch_size: batches.append(current_batch) current_batch = [] if current_batch: batches.append(current_batch) # 进程池并行处理批量数据 with Pool(num_workers) as pool: results = pool.map(process_batch, batches) # 全局去重并输出 global_unique_affils = set() for res in results: global_unique_affils.update(res) with open("final_affiliations.txt", 'w', encoding='utf-8') as f: for affil in sorted(global_unique_affils): f.write(f"{affil}\n") if __name__ == "__main__": if len(sys.argv) != 2: print("使用方式: python process_affiliations_stream.py <输入JSON数组文件路径>") sys.exit(1) # 先安装ijson: pip install ijson main(sys.argv[1])
内容的提问来源于stack exchange,提问作者JohnJ
相关产品推荐
相关产品推荐

