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

如何拆分大型JSONL文件分配至多进程,筛选德国境内闪电记录

多进程拆分方案实现

你提到的16进程拆分方案可以按以下逻辑实现,配合前置优化后处理速度可以提升1000倍以上:

  • 先做经纬度粗筛:德国的经纬度范围为纬度47.0°N ~ 55.0°N,经度5.0°E ~ 15.0°E,不在这个范围的闪电直接排除,不需要走逆地理编码,这一步可以过滤掉99%以上的无效数据,大幅减少请求量。
  • 按行分批读取大文件:不要一次性加载14GiB文件到内存,每次读取固定行数(比如2000行)作为一个任务分发给子进程,降低内存占用。
  • 子进程内批量调用逆地理编码:reverse_geocoder的search方法支持传入坐标列表做批量查询,性能比单次循环查询高2个数量级,完全避免单次查询1.5秒的问题。
  • 多进程写结果处理:可以每个子进程返回过滤后的JSONL字符串,主进程统一写结果文件,避免多进程同时写文件出现内容错乱。

示例实现代码如下:

import json
import multiprocessing as mp
import reverse_geocoder as rg

# 德国粗筛经纬度范围
DE_LAT_MIN, DE_LAT_MAX = 47.0, 55.0
DE_LON_MIN, DE_LON_MAX = 5.0, 15.0
WORKER_NUM = 16
BATCH_SIZE = 2000

def process_batch(batch_lines):
    coords = []
    stroke_refs = []
    # 先收集所有符合粗筛条件的坐标
    for line_idx, line in enumerate(batch_lines):
        record = json.loads(line)
        for stroke_idx, stroke in enumerate(record["strokes"]):
            lat, lon = stroke["lat"], stroke["lon"]
            if DE_LAT_MIN <= lat <= DE_LAT_MAX and DE_LON_MIN <= lon <= DE_LON_MAX:
                coords.append((lat, lon))
                stroke_refs.append((line_idx, stroke_idx))
    # 批量逆地理查询
    geo_results = rg.search(coords)
    # 筛选德国境内的闪电
    de_records = []
    for (line_idx, stroke_idx), res in zip(stroke_refs, geo_results):
        if res["cc"] == "DE":
            original_record = json.loads(batch_lines[line_idx])
            original_record["strokes"] = [original_record["strokes"][stroke_idx]]
            de_records.append(json.dumps(original_record))
    return "\n".join(de_records)

if __name__ == "__main__":
    pool = mp.Pool(WORKER_NUM)
    with open("input.jsonl", "r", encoding="utf-8") as f_in, open("output_de.jsonl", "w", encoding="utf-8") as f_out:
        batch = []
        for line in f_in:
            line = line.strip()
            if not line:
                continue
            batch.append(line)
            if len(batch) >= BATCH_SIZE:
                result = pool.apply_async(process_batch, args=(batch,)).get()
                if result:
                    f_out.write(result + "\n")
                batch = []
        # 处理剩余批次
        if batch:
            result = pool.apply_async(process_batch, args=(batch,)).get()
            if result:
                f_out.write(result + "\n")
    pool.close()
    pool.join()
更高性能的替代方案

如果你追求更快的处理速度,可以直接替换逆地理编码逻辑,不需要使用reverse_geocoder:

  • 用德国边界多边形直接做坐标包含判断:通过geopandas加载内置的全球国家边界数据集,提取德国的边界多边形,用shapely的contains方法直接判断坐标是否在德国境内,单进程每秒就可以处理上万次判断,比逆地理编码的性能高3个数量级以上。
    示例代码如下:
import json
import geopandas as gpd
from shapely.geometry import Point

# 加载德国边界,数据为geopandas内置,无需额外下载
world = gpd.read_file(gpd.datasets.get_path("naturalearth_lowres"))
de_boundary = world[world.iso_a3 == "DEU"].geometry.iloc[0]
# 粗筛范围和上文一致
DE_LAT_MIN, DE_LAT_MAX = 47.0, 55.0
DE_LON_MIN, DE_LON_MAX = 5.0, 15.0

def is_in_de(lat, lon):
    # 先粗筛再做精确边界判断
    if not (DE_LAT_MIN <= lat <= DE_LAT_MAX and DE_LON_MIN <= lon <= DE_LON_MAX):
        return False
    return de_boundary.contains(Point(lon, lat))

# 后续过滤逻辑直接调用is_in_de即可,不需要再调用rg.search

用这个方案即使不用多进程,也可以在十几分钟内处理完1200万行的全部数据,配合多进程处理可以把时间压缩到1分钟以内。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 23:15:00