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

百万级嵌套JSON高效解析:提取去重无邮箱机构信息并行方案咨询

百万级嵌套JSON机构信息并行解析方案(去重+移除邮箱)

需求明确

需要处理约100万条嵌套JSON记录,提取其中的Affiliation字段内容:

  1. 移除字段中的电子邮件地址(格式为Electronic address: xxx@xxx.xx)
  2. 最终输出全局唯一的机构文本条目
  3. 要求用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 05:10:56