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

如何用Python 3+实现Elasticsearch 8+批量导入JSON数据(无需CURL)

用Python 3+操作Elasticsearch 8+:创建索引并批量导入10万条JSON数据

前置准备

先安装适配Elasticsearch 8.x的Python客户端:

pip install elasticsearch

1. 连接Elasticsearch实例

Elasticsearch 8.x默认开启SSL和身份验证,根据环境选择连接方式:

  • 开发环境(仅测试用,禁用SSL):
from elasticsearch import Elasticsearch

# 连接本地ES实例,默认端口9200
es = Elasticsearch(
    "http://localhost:9200",
    basic_auth=("elastic", "你的elastic用户密码"),
    verify_certs=False
)

# 验证连接状态
if es.ping():
    print("ES连接成功")
else:
    print("ES连接失败")
  • 生产环境(启用SSL):
    需配置ES安装时生成的CA证书路径:
es = Elasticsearch(
    "https://your-es-host:9200",
    basic_auth=("elastic", "你的密码"),
    ca_certs="/path/to/http_ca.crt"
)

2. 创建自定义索引

直接通过Python API创建索引,无需CURL。先根据你的数据结构定义mapping,再执行创建操作:

from elasticsearch import Elasticsearch, exceptions

index_name = "your_data_index"

# 自定义索引mapping示例(根据实际字段调整)
index_mapping = {
    "mappings": {
        "properties": {
            "id": {"type": "integer"},
            "name": {"type": "text"},
            "created_at": {"type": "date", "format": "yyyy-MM-dd HH:mm:ss"},
            "content": {"type": "text"}
        }
    }
}

try:
    # 检查索引是否存在,不存在则创建
    if not es.indices.exists(index=index_name):
        es.indices.create(index=index_name, body=index_mapping)
        print(f"索引 {index_name} 创建成功")
    else:
        print(f"索引 {index_name} 已存在")
except exceptions.ElasticsearchException as e:
    print(f"创建索引失败: {str(e)}")

3. 批量导入10万条JSON数据

使用elasticsearch.helpers.bulk实现高效批量插入,避免单条插入的性能瓶颈:

import json
from elasticsearch.helpers import bulk

def generate_docs(json_file_path, index_name):
    """生成批量插入的文档迭代器,避免一次性加载大文件占内存"""
    with open(json_file_path, "r", encoding="utf-8") as f:
        # 情况1:JSON文件是包含所有文档的列表格式(如[{"id":1,...}, {...}])
        data = json.load(f)
        for doc in data:
            yield {
                "_index": index_name,
                "_source": doc
            }
        # 情况2:JSON文件是每行一个文档的JSON Lines格式,用下面的代码替换上面的循环
        # for line in f:
        #     doc = json.loads(line)
        #     yield {"_index": index_name, "_source": doc}

try:
    # 替换为你的JSON文件实际路径
    json_path = "/path/to/your_100k_records.json"
    # chunk_size可根据内存和ES性能调整,建议500-2000
    success, failed = bulk(es, generate_docs(json_path, index_name), chunk_size=1000)
    print(f"批量插入完成:成功 {success} 条,失败 {len(failed)} 条")
    if failed:
        print("失败的文档详情:", failed)
except Exception as e:
    print(f"批量插入失败: {str(e)}")

关键注意事项

  • 内存优化:用迭代器逐行/逐段读取JSON文件,避免一次性加载10万条数据导致内存溢出。
  • mapping适配:务必根据实际数据结构定义mapping,否则ES自动推断的字段类型可能影响后续查询。
  • 错误排查:批量操作可能存在部分文档失败,打印失败详情可快速定位问题。
  • 性能调优:chunk_size参数可根据ES集群的负载能力调整,平衡插入速度和资源占用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 20:40:44