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 文档
-
可选:清理删除整个索引,释放测试环境
更多推荐


所有评论(0)