当前位置: 首页 > news >正文

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 地址、账号、密码

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

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. 可选:清理删除整个索引,释放测试环境

http://www.jsqmd.com/news/1272722/

相关文章:

  • 计算机研究生就业数据:院校、薪资与实习的影响
  • TI TMS320TCI6484/C6457 DSP硬件设计:DDR2、JTAG与EMIF64接口实战指南
  • 深入解析TMS320DM647/DM648 DSP子系统:内存架构与性能优化实战
  • DeepSeek V4技术解析:稀疏激活与长上下文处理
  • CAD图纸自动化检测系统架构与实现
  • 【Python】常用模块:xmlrpc
  • 麒麟系统部署DeepSeek与AnythingLLM本地知识库实践
  • 嵌入式常用滤波算法与控制算法(6)互补滤波
  • 5分钟搭建免费开源的三国杀网页版:零安装的终极游戏体验
  • 怎么把管屏幕三个字拆成9张数据库表
  • C++高性能序列化与数据传输:大数据架构师的底层优化指南
  • Python自动化时间盲注脚本实战:从原理到高效渗透工具开发
  • Kimi LeetCode 3739. 统计主要元素子数组数目 II C语言实现
  • Linux网络排查利器:ss命令原理、实战与netstat替代指南
  • C++二进制文件操作:深入解析std::string序列化原理与避坑指南
  • 光伏发电系统仿真与变步长MPPT算法实践
  • 推动技术成果转化 让创新落地服务实际需求
  • Dify安装与Ollama模型接入
  • 基于Matlab的移动机器人路径规划与PID控制仿真
  • 得物App sign签名逆向:MD5加密常见错误与排查方案详解
  • 国产SPC工具在离散制造中的动态控制与边缘计算应用
  • C 语言循环与自增运算符组合对比分析
  • LangChain Memory机制详解与应用实践
  • LLM结构化输出对回答多样性的影响与平衡策略
  • 羽毛球剪辑算法集锦
  • C++异常处理终极防线:std::terminate触发机制与二次异常规避
  • 大模型技术全解析:从理论到工程实践
  • 使用JPEXS FFDec逆向分析SWF文件中的自定义加密算法与密钥生成
  • AI大模型应用开发工程师:技术落地与商业价值
  • 微信防撤回补丁失效?逆向工程实战:从内存修改到开源方案