python elasticsearch es 操作
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 地址、账号、密码 | 区分开发 / 测试 / 生产环境,敏感信息不硬编码,容器部署友好 |
|
客户端初始化连接 | Elasticsearch () 实例 + client.ping () 连通检测 | 校验 ES 服务是否正常,提前捕获连接异常,日志友好排查 |
|
模块级单例模式 | 全局变量仅初始化一次客户端,对外暴露 get_es_client () | ES 客户端内置连接池,单例复用减少 TCP 连接开销,多线程安全 | 模块加载时执行 |
简易日志封装 | log_info/log_error/log_debug 基于 print | Demo 零第三方依赖,快速查看执行结果;生产可无缝替换 logging/loguru | 区分正常日志、错误日志、调试日志 |
二、索引创建 Settings + Mappings 字段设计
知识点分类 | 核心实现 | 作用 & 生产优势 | 关键配置说明 |
索引基础 Settings | number_of_shards、number_of_replicas、refresh_interval | 分片:单机测试设 1;副本:单机 0、集群≥1;刷新间隔控制实时性 |
|
text+keyword 复合多字段 | 字符串主字段 text,内嵌 keyword 子字段 | 一套字段同时支持全文检索和精确匹配 / 排序 / 聚合,业务最通用方案 |
|
keyword 类型字段 | department_code、manager、status | 不分词、完整字符串存储,用于精确查询、分组、排序,不能全文搜索 | 编码、状态、标签、唯一标识一律用 keyword |
integer 数字类型 | employee_count | 存储数值,支持 range 范围筛选、数值排序、聚合统计 | 人数、金额、数量等数值字段 |
date 时间类型 | created_at、updated_at | 存储 ISO 标准时间字符串,支持时间区间 range 查询、时间排序 |
|
索引存在性判断 | es_client.indices.exists(index=INDEX_NAME) | 避免重复创建索引报错,幂等初始化 | 初始化索引前先判断,存在直接返回 |
删除索引 API | es_client.indices.delete(index=INDEX_NAME, ignore=[404]) | 清空全量数据,忽略索引不存在 404 报错,安全清理环境 |
|
三、单文档 CRUD 基础操作
操作类型 | ES API | 适用场景 | 核心细节 & 避坑点 |
创建单文档 | client.index() | 新增单条数据,自定义文档 ID |
|
根据 ID 精准查询 | client.get() | 根据唯一 ID 获取单条完整文档 | 文档不存在捕获异常返回 None,返回数据拼接 |
局部更新文档 | client.update (body={"doc": 更新字段}) | 仅更新传入字段,其余字段保留,推荐业务更新方式 | 只传需要修改的字段,自动刷新 |
删除单文档 | 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 数组每条包含 |
批量删除文档 | helpers.bulk + _op_type: delete | 批量根据 ID 删除数据,统一捕获失败文档,不会单条失败中断整体执行 | action 结构: |
helpers.bulk 核心特性 | success、errors 双返回值 | 不会因为个别文档失败抛出异常,可单独打印失败详情排查问题 | 返回 |
补充: 完整执行流程速查表
读取环境变量,创建全局单例 ES 客户端
初始化索引(不存在则创建,包含 settings+mappings)
单条写入多条测试部门文档
演示全部单查询场景:ID 查询、全文 match、精确 term、状态过滤、数字范围、时间范围
演示单文档局部更新、单文档删除
helpers.bulk 批量新增多条文档
helpers.bulk 批量删除指定 ID 文档
可选:清理删除整个索引,释放测试环境