python elasticsearch es 速查操作

涵盖:
  1. ES 客户端创建(环境变量配置 + 单例)
  2. 索引创建(settings + mappings,含 text/keyword 多字段)
  3. 插入文档(client.index)
  4. 按 ID 查询(client.get)
  5. 条件搜索(term / match / bool / range)
  6. 更新文档(client.update)
  7. 删除文档(client.delete)
  8. 删除索引(client.indices.delete)
  9. 批量插入、批量删除
"""
"""
Elasticsearch

涵盖:
  1. ES 客户端创建(环境变量配置 + 单例)
  2. 索引创建(settings + mappings,含 text/keyword 多字段)
  3. 插入文档(client.index)
  4. 按 ID 查询(client.get)
  5. 条件搜索(term / match / bool / range)
  6. 更新文档(client.update)
  7. 删除文档(client.delete)
  8. 删除索引(client.indices.delete)
"""

import os
from datetime import datetime
from typing import Optional, List, Dict, Any

from elasticsearch import Elasticsearch, helpers

"""用 print 代替 loguru,保持 demo 零依赖"""
def log_info(msg): print(f"[INFO] {msg}")
def log_error(msg): print(f"[ERROR] {msg}")
def log_debug(msg): pass

"""
# ===========================
# 1. ES 客户端配置
# ===========================
"""

ES_URL = os.getenv("es_url", "http://127.0.0.1:9200")
ES_USER = os.getenv("es_user", "elastic")
ES_PASSWORD = os.getenv("es_password", "elastic")

"""
索引名称
"""
INDEX_NAME = "t_department"


def create_es_client():
    """创建 Elasticsearch 客户端"""
    try:
        client = Elasticsearch(
            hosts=[ES_URL],
            basic_auth=(ES_USER, ES_PASSWORD) if ES_USER else None,
            verify_certs=False,
            request_timeout=30,
        )
        if client.ping():
            log_info(f"Elasticsearch 连接成功: {ES_URL}")
            return client
        else:
            log_error(f"Elasticsearch 连接失败: {ES_URL}")
            return None
    except Exception as e:
        log_error(f"创建 Elasticsearch 客户端失败: {e}")
        return None


"""
# 全局 ES 客户端实例(模块级单例)
"""
es_client = create_es_client()


def get_es_client():
    """获取 ES 客户端"""
    return es_client


"""
# ===========================
# 2. 索引管理
# ===========================
"""
def init_index():
    """
    初始化 t_department 索引

    字段说明:
      - department_name: text + keyword 多字段,既支持全文搜索也支持精确匹配/排序
      - department_code: keyword,部门编码,精确匹配
      - manager: keyword,部门负责人
      - employee_count: integer,员工人数
      - description: text,部门描述,全文搜索
      - status: keyword,部门状态(active / inactive)
      - created_at: date,创建时间
      - updated_at: date,更新时间
    """
    if not es_client:
        log_error("ES 客户端未初始化")
        return False

    try:
        # 检查索引是否已存在
        if es_client.indices.exists(index=INDEX_NAME):
            log_info(f"索引已存在: {INDEX_NAME}")
            return True

        # 创建索引
        es_client.indices.create(
            index=INDEX_NAME,
            body={
                "settings": {
                    "number_of_shards": 1,
                    "number_of_replicas": 0,
                    "refresh_interval": "1s",
                },
                "mappings": {
                    "properties": {
                        "department_name": {
                            "type": "text",
                            "fields": {
                                "keyword": {"type": "keyword"}
                            },
                        },
                        "department_code": {"type": "keyword"},
                        "manager": {"type": "keyword"},
                        "employee_count": {"type": "integer"},
                        "description": {"type": "text"},
                        "status": {"type": "keyword"},
                        "created_at": {"type": "date"},
                        "updated_at": {"type": "date"},
                    }
                },
            },
        )
        log_info(f"创建索引成功: {INDEX_NAME}")
        return True
    except Exception as e:
        log_error(f"初始化索引失败: {e}")
        return False
"""
# ===========================
# 3. CRUD 操作
# ===========================
"""

def create_department(
    doc_id: str,
    department_name: str,
    department_code: str,
    manager: str = "",
    employee_count: int = 0,
    description: str = "",
    status: str = "active",
) -> bool:
    """创建部门文档"""
    try:
        client = get_es_client()
        if not client:
            return False

        now = datetime.now().isoformat()
        doc = {
            "department_name": department_name,
            "department_code": department_code,
            "manager": manager,
            "employee_count": employee_count,
            "description": description,
            "status": status,
            "created_at": now,
            "updated_at": now,
        }

        client.index(index=INDEX_NAME, id=doc_id, body=doc, refresh=True)
        log_info(f"创建部门成功: {department_name} (id={doc_id})")
        return True
    except Exception as e:
        log_error(f"创建部门失败: {e}")
        return False


def get_department(doc_id: str) -> Optional[Dict[str, Any]]:
    """按 ID 获取部门"""
    try:
        client = get_es_client()
        if not client:
            return None

        result = client.get(index=INDEX_NAME, id=doc_id)
        doc = result["_source"]
        doc["id"] = result["_id"]
        return doc
    except Exception as e:
        log_debug(f"获取部门失败: {doc_id}, {e}")
        return None


def search_department_by_name(name: str) -> List[Dict[str, Any]]:
    """按部门名称全文搜索(match 查询,使用 text 字段)"""
    try:
        client = get_es_client()
        if not client:
            return []

        result = client.search(
            index=INDEX_NAME,
            body={
                "query": {"match": {"department_name": name}},
                "sort": [{"created_at": {"order": "desc"}}],
                "size": 10,
            },
        )

        docs = []
        for hit in result["hits"]["hits"]:
            doc = hit["_source"]
            doc["id"] = hit["_id"]
            docs.append(doc)
        return docs
    except Exception as e:
        log_error(f"搜索部门失败: {e}")
        return []


def search_department_by_code(code: str) -> Optional[Dict[str, Any]]:
    """按部门编码精确查询(term 查询,使用 .keyword 字段)"""
    try:
        client = get_es_client()
        if not client:
            return None

        result = client.search(
            index=INDEX_NAME,
            body={
                "query": {"term": {"department_code": code}},
                "size": 1,
            },
        )

        hits = result["hits"]["hits"]
        if hits:
            doc = hits[0]["_source"]
            doc["id"] = hits[0]["_id"]
            return doc
        return None
    except Exception as e:
        log_error(f"按编码查询部门失败: {e}")
        return None


def search_departments_by_status(status: str, limit: int = 10) -> List[Dict[str, Any]]:
    """按状态查询部门列表,按创建时间倒序"""
    try:
        client = get_es_client()
        if not client:
            return []

        result = client.search(
            index=INDEX_NAME,
            body={
                "query": {"term": {"status": status}},
                "sort": [{"created_at": {"order": "desc"}}],
                "size": limit,
            },
        )

        docs = []
        for hit in result["hits"]["hits"]:
            doc = hit["_source"]
            doc["id"] = hit["_id"]
            docs.append(doc)
        return docs
    except Exception as e:
        log_error(f"按状态查询部门失败: {e}")
        return []


def search_departments_by_employee_count(min_count: int) -> List[Dict[str, Any]]:
    """按员工人数范围查询(range 查询)"""
    try:
        client = get_es_client()
        if not client:
            return []

        result = client.search(
            index=INDEX_NAME,
            body={
                "query": {"range": {"employee_count": {"gte": min_count}}},
                "sort": [{"employee_count": {"order": "desc"}}],
                "size": 10,
            },
        )

        docs = []
        for hit in result["hits"]["hits"]:
            doc = hit["_source"]
            doc["id"] = hit["_id"]
            docs.append(doc)
        return docs
    except Exception as e:
        log_error(f"按员工人数查询部门失败: {e}")
        return []


def search_departments_by_time_range(
    start: str, end: str, limit: int = 10
) -> List[Dict[str, Any]]:
    """按创建时间范围查询"""
    try:
        client = get_es_client()
        if not client:
            return []

        result = client.search(
            index=INDEX_NAME,
            body={
                "query": {
                    "range": {
                        "created_at": {
                            "gte": start,
                            "lte": end,
                        }
                    }
                },
                "sort": [{"created_at": {"order": "desc"}}],
                "size": limit,
            },
        )

        docs = []
        for hit in result["hits"]["hits"]:
            doc = hit["_source"]
            doc["id"] = hit["_id"]
            docs.append(doc)
        return docs
    except Exception as e:
        log_error(f"按时间范围查询部门失败: {e}")
        return []


"""
是局部更新,只更新你传入的字段,其他字段保持不变。
client.update(
            index=INDEX_NAME,
            id=doc_id,
            body={"doc": update_doc},
            refresh=True,
        )
        
对应的 全量覆盖 是 client.index():
client.index(
            index=INDEX_NAME,
            id=doc_id,
            body={"doc": update_doc},
            refresh=True,
        )
"""
def update_department(
    doc_id: str,
    manager: str = None,
    employee_count: int = None,
    description: str = None,
    status: str = None,
) -> bool:
    """更新部门字段(局部更新,只更新传入的字段)"""
    try:
        client = get_es_client()
        if not client:
            return False

        update_doc = {}
        if manager is not None:
            update_doc["manager"] = manager
        if employee_count is not None:
            update_doc["employee_count"] = employee_count
        if description is not None:
            update_doc["description"] = description
        if status is not None:
            update_doc["status"] = status

        if not update_doc:
            return True

        update_doc["updated_at"] = datetime.now().isoformat()

        client.update(
            index=INDEX_NAME,
            id=doc_id,
            body={"doc": update_doc},
            refresh=True,
        )
        log_info(f"更新部门成功: {doc_id}")
        return True
    except Exception as e:
        log_error(f"更新部门失败: {doc_id}, {e}")
        return False


def delete_department(doc_id: str) -> bool:
    """删除部门"""
    try:
        client = get_es_client()
        if not client:
            return False

        client.delete(index=INDEX_NAME, id=doc_id, refresh=True)
        log_info(f"删除部门成功: {doc_id}")
        return True
    except Exception as e:
        log_error(f"删除部门失败: {doc_id}, {e}")
        return False


def delete_index():
    """删除整个索引(清空所有数据)"""
    try:
        client = get_es_client()
        if not client:
            return False

        client.indices.delete(index=INDEX_NAME, ignore=[404])
        log_info(f"删除索引成功: {INDEX_NAME}")
        return True
    except Exception as e:
        log_error(f"删除索引失败: {e}")
        return False


"""
# ===========================
# 4. Bulk 批量操作
# ===========================
"""
def bulk_create_departments(
    departments: List[Dict[str, Any]],
) -> bool:
    """
    批量创建部门文档(使用 helpers.bulk)

    参数:
        departments: 文档列表,每项格式:
            {
                "id": "dept_005",
                "department_name": "...",
                "department_code": "...",
                ... 其他字段同 create_department
            }
    """
    try:
        client = get_es_client()
        if not client:
            return False

        now = datetime.now().isoformat()
        actions = []
        for dept in departments:
            action = {
                "_index": INDEX_NAME,
                "_id": dept["id"],
                "_source": {
                    "department_name": dept["department_name"],
                    "department_code": dept["department_code"],
                    "manager": dept.get("manager", ""),
                    "employee_count": dept.get("employee_count", 0),
                    "description": dept.get("description", ""),
                    "status": dept.get("status", "active"),
                    "created_at": now,
                    "updated_at": now,
                },
            }
            actions.append(action)

        success, errors = helpers.bulk(client, actions, refresh=True)
        log_info(f"批量创建成功: {success} 条, 失败: {len(errors)} 条")
        return len(errors) == 0
    except Exception as e:
        log_error(f"批量创建失败: {e}")
        return False


def bulk_delete_departments(doc_ids: List[str]) -> bool:
    """批量删除部门文档"""
    try:
        client = get_es_client()
        if not client:
            return False

        actions = [
            {"_op_type": "delete", "_index": INDEX_NAME, "_id": doc_id} for doc_id in doc_ids
        ]

        success, errors = helpers.bulk(client, actions, refresh=True)
        log_info(f"批量删除成功: {success} 条, 失败: {len(errors)} 条")
        return len(errors) == 0
    except Exception as e:
        log_error(f"批量删除失败: {e}")
        return False


"""
# ===========================
# 5. 主程序:演示所有功能
# ===========================
"""


def main():
    print("=" * 60)
    print("Elasticsearch Demo — 基于 docparser_core 模式")
    print("=" * 60)

    if not es_client:
        print("[错误] ES 客户端未连接,请检查 ES_URL 配置")
        return

    delete_index()

    # 1. 初始化索引
    print("\n--- 1. 初始化索引 ---")
    init_index()

    # 2. 插入部门文档
    print("\n--- 2. 插入部门文档 ---")
    create_department(
        doc_id="dept_001",
        department_name="技术研发部",
        department_code="TECH",
        manager="张三",
        employee_count=50,
        description="负责公司核心产品的技术研发与架构设计",
        status="active",
    )
    create_department(
        doc_id="dept_002",
        department_name="市场营销部",
        department_code="MKT",
        manager="李四",
        employee_count=30,
        description="负责市场推广、品牌建设和销售转化",
        status="active",
    )
    create_department(
        doc_id="dept_003",
        department_name="人力资源部",
        department_code="HR",
        manager="王五",
        employee_count=15,
        description="负责招聘、培训、绩效管理和员工关系",
        status="active",
    )
    create_department(
        doc_id="dept_004",
        department_name="财务部",
        department_code="FIN",
        manager="赵六",
        employee_count=12,
        description="负责预算管理、财务报表和风险控制",
        status="inactive",
    )

    # 3. 按 ID 查询
    print("\n--- 3. 按 ID 查询 ---")
    dept = get_department("dept_001")
    if dept:
        print(f"  部门: {dept['department_name']}, 负责人: {dept['manager']}, 人数: {dept['employee_count']}")

    # 4. 全文搜索(text 字段)
    print("\n--- 4. 全文搜索(match 查询,text 字段) ---")
    results = search_department_by_name("技术")
    print(f"  搜索 '技术' 找到 {len(results)} 个部门:")
    for r in results:
        print(f"    - {r['department_name']} ({r['department_code']})")

    # 5. 精确匹配(keyword 字段)
    print("\n--- 5. 精确匹配(term 查询,keyword 字段) ---")
    dept = search_department_by_code("TECH")
    if dept:
        print(f"  编码 TECH: {dept['department_name']}")

    # 6. 按状态查询
    print("\n--- 6. 按状态查询 ---")
    active = search_departments_by_status("active")
    print(f"  活跃部门 ({len(active)} 个):")
    for r in active:
        print(f"    - {r['department_name']}")

    # 7. 范围查询(integer 字段)
    print("\n--- 7. 范围查询(range 查询,integer 字段) ---")
    big_depts = search_departments_by_employee_count(20)
    print(f"  人数 >= 20 的部门 ({len(big_depts)} 个):")
    for r in big_depts:
        print(f"    - {r['department_name']} ({r['employee_count']}人)")

    # 8. 时间范围查询(date 字段)
    print("\n--- 8. 时间范围查询(range 查询,date 字段) ---")
    now = datetime.now().isoformat()
    yesterday = datetime.now().isoformat()  # 演示用,实际可用昨天
    time_results = search_departments_by_time_range("2020-01-01T00:00:00", now)
    print(f"  2020年至今创建的部门: {len(time_results)} 个")

    # 9. 更新部门
    print("\n--- 9. 更新部门 ---")
    update_department("dept_001", manager="张三丰", employee_count=55)
    updated = get_department("dept_001")
    if updated:
        print(f"  更新后: 负责人={updated['manager']}, 人数={updated['employee_count']}")

    # 10. 删除部门
    print("\n--- 10. 删除部门 ---")
    delete_department("dept_004")
    print(f"  删除 dept_004 后, 全部活跃部门: {len(search_departments_by_status('active'))} 个")

    # 11. Bulk 批量创建部门
    print("\n--- 11. Bulk 批量创建部门 ---")
    bulk_depts = [
        {
            "id": "dept_005",
            "department_name": "产品部",
            "department_code": "PM",
            "manager": "孙七",
            "employee_count": 20,
            "description": "负责产品规划、需求分析和产品生命周期管理",
            "status": "active",
        },
        {
            "id": "dept_006",
            "department_name": "运维部",
            "department_code": "OPS",
            "manager": "周八",
            "employee_count": 18,
            "description": "负责服务器运维、监控告警和容灾管理",
            "status": "active",
        },
        {
            "id": "dept_007",
            "department_name": "法务部",
            "department_code": "LEGAL",
            "manager": "吴九",
            "employee_count": 8,
            "description": "负责合同审核、法律咨询和合规管理",
            "status": "inactive",
        },
    ]
    bulk_create_departments(bulk_depts)
    print(f"  批量创建后, 全部部门数: {len(search_departments_by_status('active')) + len(search_departments_by_status('inactive'))} 个")

    # 12. Bulk 批量删除部门
    print("\n--- 12. Bulk 批量删除部门 ---")
    bulk_delete_departments(["dept_005", "dept_007"])
    print(f"  批量删除后, 活跃部门: {len(search_departments_by_status('active'))} 个")

    # 13. 清理索引(可选,注释掉以避免误删)
    # print("\n--- 11. 清理索引 ---")
    # delete_index()
    # print(f"  索引 {INDEX_NAME} 已删除")

    print("\n" + "=" * 60)
    print("Demo 运行完毕!")
    print("=" * 60)


if __name__ == "__main__":
    main()

一、ES 客户端工程化配置

知识点分类

核心实现

作用 & 生产优势

关键代码 / 参数

环境变量解耦配置

os.getenv () 读取 ES 地址、账号、密码

区分开发 / 测试 / 生产环境,敏感信息不硬编码,容器部署友好

ES_URL = os.getenv("es_url", "默认地址")

客户端初始化连接

Elasticsearch () 实例 + client.ping () 连通检测

校验 ES 服务是否正常,提前捕获连接异常,日志友好排查

basic_auth账号密码鉴权;verify_certs=False内网关闭 SSL 校验;request_timeout=30超时防阻塞

模块级单例模式

全局变量仅初始化一次客户端,对外暴露 get_es_client ()

ES 客户端内置连接池,单例复用减少 TCP 连接开销,多线程安全

模块加载时执行create_es_client()生成全局es_client

简易日志封装

log_info/log_error/log_debug 基于 print

Demo 零第三方依赖,快速查看执行结果;生产可无缝替换 logging/loguru

区分正常日志、错误日志、调试日志

二、索引创建 Settings + Mappings 字段设计

知识点分类

核心实现

作用 & 生产优势

关键配置说明

索引基础 Settings

number_of_shards、number_of_replicas、refresh_interval

分片:单机测试设 1;副本:单机 0、集群≥1;刷新间隔控制实时性

"number_of_shards":1,"number_of_replicas":0,"refresh_interval":"1s"

text+keyword 复合多字段

字符串主字段 text,内嵌 keyword 子字段

一套字段同时支持全文检索精确匹配 / 排序 / 聚合,业务最通用方案

department_name: {type:text, fields:{keyword:{type:keyword}}}

keyword 类型字段

department_code、manager、status

不分词、完整字符串存储,用于精确查询、分组、排序,不能全文搜索

编码、状态、标签、唯一标识一律用 keyword

integer 数字类型

employee_count

存储数值,支持 range 范围筛选、数值排序、聚合统计

人数、金额、数量等数值字段

date 时间类型

created_at、updated_at

存储 ISO 标准时间字符串,支持时间区间 range 查询、时间排序

datetime.now().isoformat()生成标准时间格式

索引存在性判断

es_client.indices.exists(index=INDEX_NAME)

避免重复创建索引报错,幂等初始化

初始化索引前先判断,存在直接返回

删除索引 API

es_client.indices.delete(index=INDEX_NAME, ignore=[404])

清空全量数据,忽略索引不存在 404 报错,安全清理环境

ignore=[404]防止索引不存在抛出异常

三、单文档 CRUD 基础操作

操作类型

ES API

适用场景

核心细节 & 避坑点

创建单文档

client.index()

新增单条数据,自定义文档 ID

refresh=True写入后立即刷新,实时查询;自动填充创建 / 更新时间

根据 ID 精准查询

client.get()

根据唯一 ID 获取单条完整文档

文档不存在捕获异常返回 None,返回数据拼接id字段方便业务读取

局部更新文档

client.update (body={"doc": 更新字段})

仅更新传入字段,其余字段保留,推荐业务更新方式

只传需要修改的字段,自动刷新updated_at;区别于 index 全量覆盖

删除单文档

client.delete()

根据文档 ID 删除单条数据

文档不存在捕获异常,返回布尔值标识执行结果

四、常用 Query DSL 查询语法(Demo 全覆盖)

查询类型

适用字段类型

业务场景

核心特点

match 全文检索

text 类型主字段

模糊搜索、关键词全文匹配(如搜索部门名称含 “技术”)

对检索词分词,匹配包含分词的文档,自动计算相关性得分

term 精确匹配

keyword 字段 /.keyword 子字段

编码、状态、标签精准匹配(如部门编码 TECH、状态 active)

检索词不分词,必须与字段值完全一致才能命中

range 范围查询

integer、date

数字区间(人数≥20)、时间区间(2020 至今创建)

gte 大于等于、lte 小于等于,支持数字 / 日期两类字段

五、Bulk 批量操作(helpers 工具类)

批量操作

实现方式

优势

格式规范 & 注意事项

批量新增文档

helpers.bulk + _index/_id/_source 结构

单次请求写入多条数据,性能远高于循环单条 index;自动区分成功 / 失败条数

actions 数组每条包含_index索引名、_id文档 ID、_source完整文档数据

批量删除文档

helpers.bulk + _op_type: delete

批量根据 ID 删除数据,统一捕获失败文档,不会单条失败中断整体执行

action 结构:{"_op_type":"delete","_index":"索引名","_id":"文档ID"}

helpers.bulk 核心特性

success、errors 双返回值

不会因为个别文档失败抛出异常,可单独打印失败详情排查问题

返回(成功条数, 失败列表),判断len(errors)==0确定是否全部执行成功

补充: 完整执行流程速查表

  1. 读取环境变量,创建全局单例 ES 客户端

  2. 初始化索引(不存在则创建,包含 settings+mappings)

  3. 单条写入多条测试部门文档

  4. 演示全部单查询场景:ID 查询、全文 match、精确 term、状态过滤、数字范围、时间范围

  5. 演示单文档局部更新、单文档删除

  6. helpers.bulk 批量新增多条文档

  7. helpers.bulk 批量删除指定 ID 文档

  8. 可选:清理删除整个索引,释放测试环境

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐