如何用Python快速将大体积多语言文本转为指定JSON格式用于Elasticsearch
针对你这种大体积文本转JSON并导入Elasticsearch的需求,我整理了一套高效、兼容非英文字符的解决方案,分步骤给你讲清楚:
大文件转指定JSON格式并导入Elasticsearch方案
一、核心原则:流式处理避免内存溢出
总数据量400GB、单文件200MB的场景下,绝对不能把整个文件加载到内存中,必须采用逐行流式处理+分阶段分组的思路,既保证速度又避免内存爆掉。
二、Python实现:两种场景的代码方案
场景1:逐行生成独立JSON文档(适合直接批量导入ES)
如果每行对应一个Elasticsearch文档,结构要求包含key1和单元素的location_data数组,代码如下:
import json import os # 配置参数 INPUT_DIR = "/path/to/your/input/txt/files" OUTPUT_DIR = "/path/to/output/json/lines" BATCH_SIZE = 1000 # 每1000行生成一个JSON文件,平衡写入效率和文件大小 # 确保输出目录存在 os.makedirs(OUTPUT_DIR, exist_ok=True) def process_single_file(input_path): batch = [] batch_num = 1 with open(input_path, "r", encoding="utf-8") as f: for line in f: line = line.strip() if not line: continue # 拆分数据:假设key1和城市不含空格,描述可含空格 parts = line.split(maxsplit=2) if len(parts) != 3: print(f"跳过无效行:{line}") continue key1, city, desc = parts # 构建目标JSON结构 es_doc = { "key1": key1, "location_data": [ {"city": city, "description": desc} ] } batch.append(es_doc) # 达到批量大小就写入文件(JSON Lines格式,适合ES批量导入) if len(batch) >= BATCH_SIZE: output_file = os.path.join(OUTPUT_DIR, f"batch_{os.path.basename(input_path)}_{batch_num}.json") with open(output_file, "w", encoding="utf-8") as out_f: for item in batch: json.dump(item, out_f, ensure_ascii=False) out_f.write("\n") batch = [] batch_num += 1 # 处理剩余未批量的行 if batch: output_file = os.path.join(OUTPUT_DIR, f"batch_{os.path.basename(input_path)}_{batch_num}.json") with open(output_file, "w", encoding="utf-8") as out_f: for item in batch: json.dump(item, out_f, ensure_ascii=False) out_f.write("\n") # 遍历所有输入文件处理 for filename in os.listdir(INPUT_DIR): if filename.endswith(".txt"): input_path = os.path.join(INPUT_DIR, filename) print(f"正在处理文件:{input_path}") process_single_file(input_path)
场景2:按key1分组生成JSON文档(同一key1的城市-描述聚合到一个数组)
如果需要将同一key1的所有城市-描述条目聚合到一个location_data数组中,因为数据量大,我们先按key1哈希拆分到临时文件,再分组处理:
import json import os import hashlib # 配置参数 INPUT_DIR = "/path/to/your/input/txt/files" TEMP_DIR = "/path/to/temp/split_files" OUTPUT_DIR = "/path/to/output/grouped_json" HASH_BUCKETS = 100 # 拆分100个临时文件,控制单个文件大小 # 创建必要目录 os.makedirs(TEMP_DIR, exist_ok=True) os.makedirs(OUTPUT_DIR, exist_ok=True) # 第一步:将每行按key1哈希分到不同临时文件 def split_to_temp_files(): for filename in os.listdir(INPUT_DIR): if not filename.endswith(".txt"): continue input_path = os.path.join(INPUT_DIR, filename) print(f"正在拆分文件:{input_path}") with open(input_path, "r", encoding="utf-8") as f: for line in f: line = line.strip() if not line: continue parts = line.split(maxsplit=2) if len(parts) != 3: continue key1 = parts[0] # 计算哈希值分配到对应桶 hash_val = int(hashlib.md5(key1.encode("utf-8")).hexdigest(), 16) bucket_num = hash_val % HASH_BUCKETS temp_file = os.path.join(TEMP_DIR, f"bucket_{bucket_num}.txt") # 用|分隔存储,避免空格干扰 with open(temp_file, "a", encoding="utf-8") as temp_f: temp_f.write(f"{key1}|{parts[1]}|{parts[2]}\n") # 第二步:处理临时文件,按key1分组生成JSON def process_temp_files(): for temp_filename in os.listdir(TEMP_DIR): if not temp_filename.startswith("bucket_"): continue temp_path = os.path.join(TEMP_DIR, temp_filename) print(f"正在处理临时文件:{temp_path}") key_groups = {} with open(temp_path, "r", encoding="utf-8") as f: for line in f: line = line.strip() if not line: continue key1, city, desc = line.split("|", 2) if key1 not in key_groups: key_groups[key1] = [] key_groups[key1].append({"city": city, "description": desc}) # 写入最终JSON文件 output_file = os.path.join(OUTPUT_DIR, f"result_{temp_filename}.json") with open(output_file, "w", encoding="utf-8") as out_f: for key1, locations in key_groups.items(): grouped_doc = { "key1": key1, "location_data": locations } json.dump(grouped_doc, out_f, ensure_ascii=False) out_f.write("\n") # 可选:删除临时文件释放空间 os.remove(temp_path) # 执行流程 split_to_temp_files() process_temp_files()
三、非英文字符兼容关键细节
- 所有文件读写必须指定
encoding="utf-8"(如果原文件是GBK等其他编码,对应修改读取编码)。 - 使用
json.dump时设置ensure_ascii=False,让非英文字符直接输出,避免转成Unicode转义序列。
四、Elasticsearch批量导入优化建议
- 用JSON Lines格式:上面代码生成的每行一个JSON的格式,完美适配ES的
_bulk批量导入API。 - 批量提交:每次提交5000-10000条数据,减少HTTP请求次数。
- 临时关闭刷新:导入前执行
PUT /your_index/_settings {"refresh_interval": "-1"},导入完成后恢复为"refresh_interval": "30s",大幅提升导入速度。 - 多进程处理:可以用
multiprocessing模块同时处理多个输入文件,利用多核CPU加速转换。
示例批量导入ES的代码片段:
from elasticsearch import Elasticsearch from elasticsearch.helpers import bulk # 初始化ES客户端 es = Elasticsearch(["http://your-es-host:9200"], basic_auth=("username", "password")) INDEX_NAME = "your_target_index" def import_json_to_es(json_file_path): actions = [] with open(json_file_path, "r", encoding="utf-8") as f: for line in f: doc = json.loads(line) actions.append({ "_index": INDEX_NAME, "_source": doc }) if len(actions) >= 5000: bulk(es, actions) actions = [] # 处理剩余数据 if actions: bulk(es, actions) # 遍历所有输出JSON文件导入 for filename in os.listdir(OUTPUT_DIR): if filename.endswith(".json"): import_json_to_es(os.path.join(OUTPUT_DIR, filename))
内容的提问来源于stack exchange,提问作者creed
相关产品推荐
相关产品推荐

