python elasticsearch es 操作

发布时间:2026/7/27 7:52:45
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_certsFalse, request_timeout30, ) if client.ping(): log_info(fElasticsearch 连接成功: {ES_URL}) return client else: log_error(fElasticsearch 连接失败: {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(indexINDEX_NAME): log_info(f索引已存在: {INDEX_NAME}) return True # 创建索引 es_client.indices.create( indexINDEX_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(indexINDEX_NAME, iddoc_id, bodydoc, refreshTrue) 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(indexINDEX_NAME, iddoc_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( indexINDEX_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( indexINDEX_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( indexINDEX_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( indexINDEX_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( indexINDEX_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( indexINDEX_NAME, iddoc_id, body{doc: update_doc}, refreshTrue, ) 对应的 全量覆盖 是 client.index() client.index( indexINDEX_NAME, iddoc_id, body{doc: update_doc}, refreshTrue, ) 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( indexINDEX_NAME, iddoc_id, body{doc: update_doc}, refreshTrue, ) 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(indexINDEX_NAME, iddoc_id, refreshTrue) 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(indexINDEX_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, refreshTrue) 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, refreshTrue) 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_iddept_001, department_name技术研发部, department_codeTECH, manager张三, employee_count50, description负责公司核心产品的技术研发与架构设计, statusactive, ) create_department( doc_iddept_002, department_name市场营销部, department_codeMKT, manager李四, employee_count30, description负责市场推广、品牌建设和销售转化, statusactive, ) create_department( doc_iddept_003, department_name人力资源部, department_codeHR, manager王五, employee_count15, description负责招聘、培训、绩效管理和员工关系, statusactive, ) create_department( doc_iddept_004, department_name财务部, department_codeFIN, manager赵六, employee_count12, description负责预算管理、财务报表和风险控制, statusinactive, ) # 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_count55) 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_certsFalse内网关闭 SSL 校验request_timeout30超时防阻塞模块级单例模式全局变量仅初始化一次客户端对外暴露 get_es_client ()ES 客户端内置连接池单例复用减少 TCP 连接开销多线程安全模块加载时执行create_es_client()生成全局es_client简易日志封装log_info/log_error/log_debug 基于 printDemo 零第三方依赖快速查看执行结果生产可无缝替换 logging/loguru区分正常日志、错误日志、调试日志二、索引创建 Settings Mappings 字段设计知识点分类核心实现作用 生产优势关键配置说明索引基础 Settingsnumber_of_shards、number_of_replicas、refresh_interval分片单机测试设 1副本单机 0、集群≥1刷新间隔控制实时性number_of_shards:1,number_of_replicas:0,refresh_interval:1stextkeyword 复合多字段字符串主字段 text内嵌 keyword 子字段一套字段同时支持全文检索和精确匹配 / 排序 / 聚合业务最通用方案department_name: {type:text, fields:{keyword:{type:keyword}}}keyword 类型字段department_code、manager、status不分词、完整字符串存储用于精确查询、分组、排序不能全文搜索编码、状态、标签、唯一标识一律用 keywordinteger 数字类型employee_count存储数值支持 range 范围筛选、数值排序、聚合统计人数、金额、数量等数值字段date 时间类型created_at、updated_at存储 ISO 标准时间字符串支持时间区间 range 查询、时间排序datetime.now().isoformat()生成标准时间格式索引存在性判断es_client.indices.exists(indexINDEX_NAME)避免重复创建索引报错幂等初始化初始化索引前先判断存在直接返回删除索引 APIes_client.indices.delete(indexINDEX_NAME, ignore[404])清空全量数据忽略索引不存在 404 报错安全清理环境ignore[404]防止索引不存在抛出异常三、单文档 CRUD 基础操作操作类型ES API适用场景核心细节 避坑点创建单文档client.index()新增单条数据自定义文档 IDrefreshTrue写入后立即刷新实时查询自动填充创建 / 更新时间根据 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确定是否全部执行成功补充 完整执行流程速查表读取环境变量创建全局单例 ES 客户端初始化索引不存在则创建包含 settingsmappings单条写入多条测试部门文档演示全部单查询场景ID 查询、全文 match、精确 term、状态过滤、数字范围、时间范围演示单文档局部更新、单文档删除helpers.bulk 批量新增多条文档helpers.bulk 批量删除指定 ID 文档可选清理删除整个索引释放测试环境