如何拆分大型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
相关产品推荐
相关产品推荐

