职位搜索from datetime import date, datetime from enum import Enum from typing import Any, Optional from elasticsearch import AsyncElasticsearch from elasticsearch.helpers import async_bulk from fastapi import Depends, APIRouter, Query 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, }, } es_data_router.get(/search, summary职位搜索) async def search_jobs( # ---------- 关键词 ---------- keyword: Optional[str] Query(None, description搜索关键词如 Java、前端), # ---------- 筛选 ---------- work_location: Optional[str] Query(None, description工作地点精确匹配), edu_require: Optional[str] Query(None, description学历要求), exp_require: Optional[str] Query(None, description经验要求), industry_id: Optional[int] Query(None, description行业ID), company_scale: Optional[int] Query(None, description公司规模枚举值), financing_stage: Optional[int] Query(None, description融资阶段枚举值), status: Optional[int] Query(1, description职位状态默认1招聘中), # ---------- 分页 / 排序 ---------- page: int Query(1, ge1, description页码), page_size: int Query(10, ge1, le50, description每页条数), sort: str Query(time, description排序score相关度time发布时间), es_client: AsyncElasticsearch Depends(es_client_depend), ): 职位搜索思路 1. bool.must/should关键词 multi_match有 keyword 才加 2. bool.filter城市、学历、状态等精确条件不参与算分 3. from/size分页 4. sort有词按相关度或按 publish_time # 索引不存在时友好提示 if not await es_client.indices.exists(indexBOSS_JOB_INDEX_NAME_V2): return { code: 0, message: f索引 {BOSS_JOB_INDEX_NAME_V2} 不存在请先创建并同步数据, } must_clauses: list[dict[str, Any]] [] filter_clauses: list[dict[str, Any]] [] # ----- 1关键词打在多个 text 字段职位名称加权 ----- if keyword and keyword.strip(): must_clauses.append( { multi_match: { query: keyword.strip(), fields: [ job_name^3, # 名称命中权重更高 job_desc^2, duty_require, enterprise_name, industry_name, ], type: best_fields, operator: and, # 可按需改成 or召回更宽 } } ) # ----- 2硬筛选全部放 filter ----- if status is not None: filter_clauses.append({term: {status: status}}) if work_location: filter_clauses.append({term: {work_location: work_location}}) if edu_require: filter_clauses.append({term: {edu_require: edu_require}}) if exp_require: filter_clauses.append({term: {exp_require: exp_require}}) if industry_id is not None: filter_clauses.append({term: {industry_id: industry_id}}) if company_scale is not None: filter_clauses.append({term: {enterpriseInfo_company_scale: company_scale}}) if financing_stage is not None: filter_clauses.append({term: {enterpriseInfo_financing_stage: financing_stage}}) # ----- 3组装 bool ----- bool_query: dict[str, Any] {} if must_clauses: bool_query[must] must_clauses if filter_clauses: bool_query[filter] filter_clauses # 既无关键词也无筛选时匹配全部仍建议至少有默认 status filter query: dict[str, Any] {bool: bool_query} if bool_query else {match_all: {}} # ----- 4排序 ----- # 有关键词且 sortscore → 按相关度否则按发布时间倒序 if sort score and keyword: sort_clause: Any [_score, {publish_time: {order: desc, missing: _last}}] else: sort_clause [{publish_time: {order: desc, missing: _last}}] # ----- 5分页ES 用 from / size ----- from_ (page - 1) * page_size # 列表页只取卡片字段减小 _source source_fields [ job_id, job_name, work_location, min_salary, max_salary, salary_times, edu_require, exp_require, job_tags, status, publish_time, enterprise_id, enterprise_name, enterprise_city_name, enterpriseInfo_company_scale, enterpriseInfo_financing_stage, industry_id, industry_name, ] resp await es_client.search( indexBOSS_JOB_INDEX_NAME_V2, queryquery, sortsort_clause, from_from_, sizepage_size, sourcesource_fields, track_total_hitsTrue, # 拿到准确 total ) hits resp.get(hits, {}) total hits.get(total, {}) # ES 7 total 一般是 {value: n, relation: eq} total_count total.get(value, 0) if isinstance(total, dict) else int(total or 0) lists [] for hit in hits.get(hits, []): item hit.get(_source) or {} item[_score] hit.get(_score) # 可选方便调相关度 lists.append(item) return { code: 1, message: success, data: { lists: lists, page_info: { page: page, page_size: page_size, total_count: total_count, total_page: (total_count page_size - 1) // page_size if page_size else 0, }, }, }岗位详细job_router.get(/info/{job_id}, summary职位详情) async def getJobInfoById(job_id: int): res await JobService.getJobInfoById(job_id) return { code: 1, message: 查询成功, data: res } staticmethod async def getJobInfoById(job_id: int): job await Job.get_or_none(idjob_id) enterprise_id job.enterprise_id enterprise await Enterprise.get_or_none(identerprise_id) info { job_info: job, enterprise_info: enterprise } return info简历投递创建简历投递记录表class ResumeStatus(IntEnum): 简历状态 SUBMITTED 1 # 已投递 REVIEWING 2 # 以查看 class ResumeSubmissionRecords(Model): 职位投递记录表 id fields.IntField(pkTrue, description主键自增ID) job_id fields.IntField(description职位ID) job_seeker_id fields.IntField(description求职者ID) # 附件简历URL copy_resume_url fields.CharField(max_length256, description附件简历URL) # 投递时间 submit_time fields.DatetimeField(auto_now_addTrue, description投递时间) update_time fields.DatetimeField(auto_nowTrue, description更新时间) resume_status fields.IntEnumField(enum_typeResumeStatus, description简历状态) class Meta: table t_resume_submission_records table_description 职位投递记录表job_router.post(/resume_submission, summary职位投递) async def resume_submission(job_seeker_idDepends(get_job_seeker_info), job_id: int Query(..., description职位ID)): await JobService.resume_submission(job_seeker_id, job_id) return {code: 1, message: 简历投递成功} staticmethod async def resume_submission(job_seeker_id: int, job_id: int): oss AliyunOSSTool() attachmentResume await AttachmentResume.get_or_none(job_seeker_idjob_seeker_id, is_defaultTrue) file_url attachmentResume.file_url source_key file_url.replace(fhttps://{oss.bucket_name}.{oss.endpoint}/, ) logger.info(fsource_key:{source_key}) snowflake SnowflakeSingleton.get_instance(worker_id1) file_id snowflake.get_id() # 修复提取纯文件名避免路径重复 source_filename os.path.basename(source_key) # 只保留文件名如e3d73f50-入职通知书.docx target_key fjob_seeker_avatar/copy/{file_id}-{source_filename} # 正确路径 copy_ok, copy_data oss.copy_single_file( source_oss_keysource_key, target_oss_keytarget_key, overwriteFalse ) logger.info(fcopy_ok:{copy_ok}) logger.info(fcopy_data:{copy_data}) await ResumeSubmissionRecords.create( job_idjob_id, job_seeker_idjob_seeker_id, copy_resume_urlcopy_data[target_url], submit_timenow(), resume_statusResumeStatus.SUBMITTED )简历中心enterprise_router.get(/resume_center, summary简历中心, description简历中心) async def resume_center(team_idDepends(get_job_info)): res await EnterpriseService.resume_center(team_id) return { code: 1, message: 成功, data: res } staticmethod async def resume_center(team_id: int): jobs await Job.filter(recruit_team_idteam_id) data_dict_list [] for job in jobs: job_id job.id resume_submission_records await ResumeSubmissionRecords.filter(job_idjob_id) for resume_submission_record in resume_submission_records: job_seeker_id resume_submission_record.job_seeker_id job_seeker await JobSeeker.get_or_none(idjob_seeker_id) jobIntention await JobIntention.get_or_none(job_seeker_idjob_seeker_id) workExperience await WorkExperience.filter(job_seeker_idjob_seeker_id).order_by( -entry_time).first() info_dict { job_seeker_id: job_seeker_id, job_seeker: job_seeker, jobIntention: jobIntention, workExperience: workExperience, resume_status: resume_submission_record.resume_status, submit_time: resume_submission_record.submit_time, } data_dict_list.append(info_dict) return data_dict_list