练习项目跟进es创建(day09) es官网Elasticsearch Python 客户端发布说明 |蟒蛇https://www.elastic.co/docs/release-notes/elasticsearch/clients/python入门指南 |蟒蛇https://www.elastic.co/docs/reference/elasticsearch/clients/python/getting-started官网页面步骤1、安装python -m pip install elasticsearch[async]2、链接import osfrom elasticsearch import AsyncElasticsearchclient AsyncElasticsearch(https://...,api_keyos.environ[ELASTIC_API_KEY], # 是否设置了账号和密码没有设置可以忽略)3、运营4、创建索引await client.indices.create(indexmy_index)mappings {properties: {foo: {type: text},bar: {type: text,fields: {keyword: {type: keyword,ignore_above: 256,}},},}}await client.indices.create(indexmy_index, mappingsmappings)核心字段类型及适用场景1、字符串text特点会进行分词处理支持全文搜索使用场景存储长文本内容如商品描述、文章内容、评论等需要全文检索的场景示例商品的详细介绍、用户评价内容keyword特点不进行分词精确匹配支持聚合和排序使用场景存储ID、编号、标签、状态等需要精确匹配或分组统计的字段示例商品ID、订单编号、商品分类、用户标签2、数值类型integer特点32位整数范围-2³¹到2³¹-1使用场景存储数量、年龄等整数数据示例商品库存数量、用户年龄long特点64位整数范围-2⁶³到2⁶³-1使用场景存储较大的整数如订单金额分、用户ID等示例订单总金额以分为单位、用户唯一标识float特点32位单精度浮点数使用场景存储精度要求不高的小数示例商品评分1-5分double特点64位双精度浮点数使用场景存储精度要求较高的小数示例商品价格、重量等3、日期类型date特点支持多种日期格式可进行范围查询和排序使用场景存储创建时间、更新时间、过期时间等示例商品上架时间、订单创建时间、活动结束时间4、布尔类型boolean特点只能是true或false使用场景存储二值状态示例商品是否上架、订单是否支付、用户是否激活5、数组类型array特点没有专门的数组类型任何字段都可以存储多个值使用场景存储多个相同类型的值示例商品标签多个标签、用户拥有的权限多个权限6、对象类型object特点用于存储嵌套的JSON对象使用场景存储具有复杂结构的数据示例商品的规格信息包含颜色、尺寸、材质等子字段7、地理坐标类型geo_point特点存储经纬度坐标使用场景需要根据地理位置进行搜索和距离计算示例店铺位置、用户收货地址位置5、索引文档await client.index(indexmy_index,idmy_document_id,document{foo: foo,bar: bar,})6、获取文件await client.get(indexmy_index, idmy_document_id)es_client.py ES 客户端模块用于与 Elasticsearch 进行交互。 from elasticsearch import AsyncElasticsearch import os from app.core.logging import logger ES_HOST os.getenv(ES_HOST, http://localhost:9200) es_client: AsyncElasticsearch | None None async def get_es_client() - AsyncElasticsearch: 获取 ES 异步客户端单例 global es_client if es_client is None: es_client AsyncElasticsearch(ES_HOST) logger.info(ES 客户端初始化成功) return es_client async def close_es_client(): 关闭 ES 客户端连接 global es_client if es_client is not None: await es_client.close() es_client Nonedepends.pyfrom elasticsearch import AsyncElasticsearch # 依赖项, 用于注入 Elasticsearch 客户端 async def es_client_depend() - AsyncElasticsearch: return await get_es_client()es_data_api.pyfrom datetime import date, datetime from enum import Enum from typing import Any from elasticsearch import AsyncElasticsearch from elasticsearch.helpers import async_bulk from fastapi import Depends, APIRouter from app.core.depends import es_client_depend from app.core.logging import logger from app.models import Job, Enterprise, EnterpriseInfo # https://www.elastic.co/docs/reference/elasticsearch/clients/python/getting-started es_data_router APIRouter( prefix/es-data, tags[elasticsearch测试], ) BOSS_JOB_INDEX_NAME boss_job_index # 优化版接口使用独立索引名避免和旧接口互相覆盖便于课堂对比 BOSS_JOB_INDEX_NAME_V2 boss_job_index_v2 def _to_es_value(value: Any) - Any: 把 ORM 字段值转成 ES 友好的可 JSON 序列化类型。 - datetime / date → ISO 字符串ES date 字段可识别 - Enum / IntEnum → 对应 value一般为 int - 其余原样返回含 None、str、list、dict if value is None: return None if isinstance(value, datetime): return value.isoformat() if isinstance(value, date): return value.isoformat() if isinstance(value, Enum): return value.value return value def _build_job_document( job: Job, enterprise: Enterprise | None, enterprise_info: EnterpriseInfo | None, ) - dict[str, Any]: 将「职位 企业 企业详情」拼成一份扁平文档宽表。 设计要点 1. 企业/详情缺失时填 None不抛异常保证整批同步不被单条脏数据打断 2. 字段名与 create-index-v2 的 mapping 一一对应 3. 统一走 _to_es_value避免 datetime/Enum 序列化问题 city getattr(enterprise, city, None) if enterprise else None industry getattr(enterprise_info, industry, None) if enterprise_info else None return { # ---------- 职位本身 ---------- job_id: job.id, job_name: job.job_name, department_id: _to_es_value(job.department_id), work_location: job.work_location, # 模型里薪资是 CharField可能含「面议」因此 mapping 用 keyword这里保持字符串 min_salary: job.min_salary, max_salary: job.max_salary, salary_times: job.salary_times, edu_require: job.edu_require, exp_require: job.exp_require, gender_require: job.gender_require, # 模型里招聘人数也是字符串用 keyword 更稳妥 recruit_num: job.recruit_num, # JSONField常见形态是 [五险一金,年终奖]keyword 支持多值 job_tags: job.job_tags or [], job_desc: job.job_desc, duty_require: job.duty_require, status: _to_es_value(job.status), publish_time: _to_es_value(job.publish_time), enterprise_id: job.enterprise_id, recruit_team_id: job.recruit_team_id, # ---------- 企业主表可空 ---------- enterprise_name: enterprise.enterprise_name if enterprise else None, enterprise_code: enterprise.enterprise_code if enterprise else None, enterprise_city_id: city.id if city else None, enterprise_city_name: city.name if city else None, enterprise_account_status: _to_es_value(enterprise.account_status) if enterprise else None, enterprise_create_time: _to_es_value(enterprise.create_time) if enterprise else None, enterprise_auth_time: _to_es_value(enterprise.auth_time) if enterprise else None, enterprise_auth_type: _to_es_value(enterprise.auth_type) if enterprise else None, enterprise_risk_level: _to_es_value(enterprise.risk_level) if enterprise else None, enterprise_blacklist_status: _to_es_value(enterprise.blacklist_status) if enterprise else None, enterprise_complaint_count: enterprise.complaint_count if enterprise else None, enterprise_company_website: enterprise.company_website if enterprise else None, enterprise_email: enterprise.email if enterprise else None, enterprise_audit_type: _to_es_value(enterprise.audit_type) if enterprise else None, enterprise_submit_time: _to_es_value(enterprise.submit_time) if enterprise else None, # ---------- 企业详情可空 ---------- enterpriseInfo_unified_social_credit_code: ( enterprise_info.unified_social_credit_code if enterprise_info else None ), enterpriseInfo_legal_representative: ( enterprise_info.legal_representative if enterprise_info else None ), enterpriseInfo_registered_capital: ( enterprise_info.registered_capital if enterprise_info else None ), enterpriseInfo_establish_date: ( _to_es_value(enterprise_info.establish_date) if enterprise_info else None ), enterpriseInfo_register_status: ( _to_es_value(enterprise_info.register_status) if enterprise_info else None ), # 规模/融资在模型里是 IntEnum这里存 int便于精确筛选 enterpriseInfo_company_scale: ( _to_es_value(enterprise_info.company_scale) if enterprise_info else None ), enterpriseInfo_financing_stage: ( _to_es_value(enterprise_info.financing_stage) if enterprise_info else None ), enterpriseInfo_headquarters_address: ( enterprise_info.headquarters_address if enterprise_info else None ), enterpriseInfo_business_scope: ( enterprise_info.business_scope if enterprise_info else None ), # ---------- 行业 ---------- industry_id: industry.id if industry else None, industry_name: industry.name if industry else None, } es_data_router.post(/create-index, summary创建索引) async def create_index(es_client: AsyncElasticsearch Depends(es_client_depend)): mappings { properties: { job_id: { type: long }, job_name: { type: text, analyzer: ik_max_word }, department_id: { type: long }, work_location: { type: keyword }, min_salary: { type: keyword }, max_salary: { type: keyword }, salary_times: { type: keyword }, edu_require: { type: keyword, }, exp_require: { type: keyword }, gender_require: { type: keyword }, recruit_num: { type: long }, job_tags: { type: keyword }, job_desc: { type: text, analyzer: ik_max_word }, duty_require: { type: text, analyzer: ik_max_word }, status: { type: integer }, publish_time: { type: date, # format: yyyy-MM-dd HH:mm:ss }, enterprise_id: { type: long }, recruit_team_id: { type: long }, enterprise_name: { type: text, analyzer: ik_max_word }, enterprise_code: { type: keyword }, enterprise_city_id: { type: long }, enterprise_city_name: { type: keyword }, enterprise_account_status: { type: integer }, enterprise_create_time: { type: date, # format: yyyy-MM-dd HH:mm:ss }, enterprise_auth_time: { type: date, # format: yyyy-MM-dd HH:mm:ss }, enterprise_auth_type: { type: integer }, enterprise_risk_level: { type: integer }, enterprise_blacklist_status: { type: integer }, enterprise_complaint_count: { type: integer }, enterprise_company_website: { type: keyword }, enterprise_email: { type: keyword }, enterprise_audit_type: { type: integer }, enterprise_submit_time: { type: date, # format: yyyy-MM-dd HH:mm:ss }, enterpriseInfo_unified_social_credit_code: { type: keyword }, enterpriseInfo_legal_representative: { type: keyword }, enterpriseInfo_registered_capital: { type: keyword}, enterpriseInfo_establish_date: { type: date, # format: yyyy-MM-dd HH:mm:ss }, enterpriseInfo_register_status: { type: integer }, enterpriseInfo_company_scale: { type: keyword }, enterpriseInfo_financing_stage: { type: keyword }, enterpriseInfo_headquarters_address: { type: keyword }, industry_id: { type: long }, industry_name: { type: text, analyzer: ik_max_word } } } if await es_client.indices.exists(indexBOSS_JOB_INDEX_NAME): return 索引已存在 await es_client.indices.create( indexBOSS_JOB_INDEX_NAME, mappingsmappings ) return 索引创建成功 es_data_router.post(/insert-data, summary同步数据) async def insert_data(es_client: AsyncElasticsearch Depends(es_client_depend)): jobs await Job.all() for job in jobs: enterprise_id job.enterprise_id enterprise await Enterprise.get_or_none(identerprise_id).prefetch_related(city) enterpriseInfo await EnterpriseInfo.get_or_none(enterprise_identerprise_id).prefetch_related(industry) job_info { job_id: job.id, job_name: job.job_name, department_id: job.department_id, work_location: job.work_location, min_salary: job.min_salary, max_salary: job.max_salary, salary_times: job.salary_times, edu_require: job.edu_require, exp_require: job.exp_require, gender_require: job.gender_require, recruit_num: job.recruit_num, job_tags: job.job_tags, job_desc: job.job_desc, duty_require: job.duty_require, status: job.status, publish_time: job.publish_time, enterprise_id: job.enterprise_id, recruit_team_id: job.recruit_team_id, enterprise_name: enterprise.enterprise_name, enterprise_code: enterprise.enterprise_code, # enterprise_city_id: enterprise.city.id, # enterprise_city_name: job.work_location, enterprise_account_status: enterprise.account_status, enterprise_create_time: enterprise.create_time, enterprise_auth_time: enterprise.auth_time, enterprise_auth_type: enterprise.auth_type, enterprise_risk_level: enterprise.risk_level, enterprise_blacklist_status: enterprise.blacklist_status, enterprise_complaint_count: enterprise.complaint_count, enterprise_company_website: enterprise.company_website, enterprise_email: enterprise.email, enterprise_audit_type: enterprise.audit_type, enterprise_submit_time: enterprise.submit_time, enterpriseInfo_unified_social_credit_code: enterpriseInfo.unified_social_credit_code, enterpriseInfo_legal_representative: enterpriseInfo.legal_representative, enterpriseInfo_registered_capital: enterpriseInfo.registered_capital, enterpriseInfo_establish_date: enterpriseInfo.establish_date, enterpriseInfo_register_status: enterpriseInfo.register_status, enterpriseInfo_company_scale: enterpriseInfo.company_scale, enterpriseInfo_financing_stage: enterpriseInfo.financing_stage, enterpriseInfo_headquarters_address: enterpriseInfo.headquarters_address, enterpriseInfo_business_scope: enterpriseInfo.business_scope, industry_id: enterpriseInfo.industry.id, industry_name: enterpriseInfo.industry.name, } await es_client.index( indexBOSS_JOB_INDEX_NAME, documentjob_info ) return 数据同步成功 # # 优化版接口保留上面旧接口不动便于对比学习 # 主要改进 # 1. mapping 与写入字段对齐补 business_scope、城市字段枚举用 integer # 2. 同步时指定文档 _idjob.id重复调用可幂等覆盖不会越插越多 # 3. 批量预加载企业/详情避免 N1 # 4. 企业缺失不抛错跳过或写空字段并记录日志 # 5. datetime/Enum 统一序列化使用 async_bulk 批量写入 # es_data_router.post(/create-index-v2, summary创建索引(优化版)) async def create_index_v2(es_client: AsyncElasticsearch Depends(es_client_depend)): 创建优化版职位搜索索引 boss_job_index_v2。 与旧版 /create-index 的区别 - 使用独立索引名不覆盖旧索引 - mapping 补齐 enterpriseInfo_business_scope、城市相关字段 - company_scale / financing_stage / register_status 用 integer与 IntEnum 一致 - recruit_num 改为 keyword与模型 CharField 一致避免「若干」等非数字写不进去 - 需要 ES 已安装 IK 分词插件ik_max_word # settings可按需加分片/副本单机 Docker 开发一般 1 分片 0 副本即可 settings { number_of_shards: 1, number_of_replicas: 0, } # mappings定义每个字段如何被索引与查询 # - text ik_max_word中文全文检索 # - keyword精确匹配、聚合、筛选 # - integer/long数值筛选 # - date时间范围查询写入时用 ISO 字符串 mappings { properties: { # ----- 职位 ----- job_id: {type: long}, job_name: {type: text, analyzer: ik_max_word}, department_id: {type: integer}, work_location: {type: keyword}, min_salary: {type: keyword}, max_salary: {type: keyword}, salary_times: {type: keyword}, edu_require: {type: keyword}, exp_require: {type: keyword}, gender_require: {type: keyword}, recruit_num: {type: keyword}, job_tags: {type: keyword}, job_desc: {type: text, analyzer: ik_max_word}, duty_require: {type: text, analyzer: ik_max_word}, status: {type: integer}, publish_time: {type: date}, enterprise_id: {type: long}, recruit_team_id: {type: long}, # ----- 企业主表 ----- enterprise_name: {type: text, analyzer: ik_max_word}, enterprise_code: {type: keyword}, enterprise_city_id: {type: long}, enterprise_city_name: {type: keyword}, enterprise_account_status: {type: integer}, enterprise_create_time: {type: date}, enterprise_auth_time: {type: date}, enterprise_auth_type: {type: integer}, enterprise_risk_level: {type: integer}, enterprise_blacklist_status: {type: integer}, enterprise_complaint_count: {type: integer}, enterprise_company_website: {type: keyword}, enterprise_email: {type: keyword}, enterprise_audit_type: {type: integer}, enterprise_submit_time: {type: date}, # ----- 企业详情 ----- enterpriseInfo_unified_social_credit_code: {type: keyword}, enterpriseInfo_legal_representative: {type: keyword}, enterpriseInfo_registered_capital: {type: keyword}, enterpriseInfo_establish_date: {type: date}, enterpriseInfo_register_status: {type: integer}, enterpriseInfo_company_scale: {type: integer}, enterpriseInfo_financing_stage: {type: integer}, enterpriseInfo_headquarters_address: {type: keyword}, # 旧版 mapping 漏了该字段同步却在写 → 这里显式声明 enterpriseInfo_business_scope: { type: text, analyzer: ik_max_word, }, # ----- 行业 ----- industry_id: {type: long}, industry_name: {type: text, analyzer: ik_max_word}, } } # 已存在则不重复创建需要重建可手动 DELETE 索引后再调本接口 if await es_client.indices.exists(indexBOSS_JOB_INDEX_NAME_V2): return { code: 1, message: 索引已存在, data: {index: BOSS_JOB_INDEX_NAME_V2}, } await es_client.indices.create( indexBOSS_JOB_INDEX_NAME_V2, settingssettings, mappingsmappings, ) return { code: 1, message: 索引创建成功, data: {index: BOSS_JOB_INDEX_NAME_V2}, } es_data_router.post(/insert-data-v2, summary同步数据(优化版)) async def insert_data_v2(es_client: AsyncElasticsearch Depends(es_client_depend)): 全量同步职位数据到 boss_job_index_v2。 与旧版 /insert-data 的区别 1. 文档 _id 使用 job.id → 重复同步会覆盖不会产生重复文档 2. 先批量查出 Enterprise / EnterpriseInfo再内存关联 → 避免 N1 3. 企业或详情缺失时记录日志并跳过该职位不中断整批 4. 使用 async_bulk 批量写入比逐条 index 更快 5. 日期/枚举统一序列化后再写入 使用前请先调用 /create-index-v2。 # 索引不存在时提前提示避免 bulk 时才报错 if not await es_client.indices.exists(indexBOSS_JOB_INDEX_NAME_V2): return { code: 0, message: f索引 {BOSS_JOB_INDEX_NAME_V2} 不存在请先调用 /es-data/create-index-v2, } # 1一次取出全部职位 jobs await Job.all() if not jobs: return {code: 1, message: 没有可同步的职位, data: {success: 0, skip: 0}} # 2收集涉及到的企业 ID批量查询企业与详情解决 N1 enterprise_ids list({job.enterprise_id for job in jobs if job.enterprise_id is not None}) enterprises await Enterprise.filter(id__inenterprise_ids).prefetch_related(city) enterprise_map {e.id: e for e in enterprises} enterprise_infos await EnterpriseInfo.filter( enterprise_id__inenterprise_ids ).prefetch_related(industry) # EnterpriseInfo 按企业 ID 建索引方便 O(1) 查找 enterprise_info_map {info.enterprise_id: info for info in enterprise_infos} # 3组装 bulk 动作列表 # async_bulk 接受的每个 action 至少包含_index、_id可选但强烈建议、_source 或直接字段 actions: list[dict[str, Any]] [] skip_count 0 for job in jobs: enterprise enterprise_map.get(job.enterprise_id) enterprise_info enterprise_info_map.get(job.enterprise_id) # 没有企业主数据时跳过宽表缺核心信息写入后搜索意义不大 if enterprise is None: skip_count 1 logger.warning(f同步跳过职位 id{job.id} 关联企业 id{job.enterprise_id} 不存在) continue # 详情缺失允许继续写企业主表字段仍可用industry 等会是 None if enterprise_info is None: logger.warning(f职位 id{job.id} 无企业详情将写入空的 enterpriseInfo_* 字段) document _build_job_document(job, enterprise, enterprise_info) # _idstr(job.id)幂等的关键。再次同步同一职位会覆盖旧文档而不是新增一条 actions.append( { _index: BOSS_JOB_INDEX_NAME_V2, _id: str(job.id), _source: document, } ) if not actions: return { code: 1, message: 没有成功组装的文档可能全部被跳过, data: {success: 0, skip: skip_count}, } # 4批量写入raise_on_errorFalse 时单条失败不会整体抛异常由返回值统计 success_count, errors await async_bulk( clientes_client, actionsactions, raise_on_errorFalse, ) # errors 在 raise_on_errorFalse 时是失败详情列表 error_count len(errors) if isinstance(errors, list) else 0 if error_count: logger.error(fES bulk 部分失败失败条数{error_count}样例{errors[:3]}) return { code: 1, message: 数据同步完成, data: { index: BOSS_JOB_INDEX_NAME_V2, job_total: len(jobs), success: success_count, skip: skip_count, error: error_count, }, }

本月热点