FastAPI + Tortoise-ORM + Elasticsearch 实战:职位数据全量同步到 ES 的设计
项目实践FastAPI Elasticsearch 职位数据全量同步FastAPI Tortoise-ORM Elasticsearch 实战职位数据全量同步到 ES 的设计意识与踩坑复盘一、前言介绍1.1 功能定位1.2 数据模型总览1.3 同步流程总览二、环境准备2.1 依赖与 ES 客户端2.2 索引与运行环境2.3 路由与生命周期挂载三、知识点讲解3.1 ORM 与 ES 的数据形态差异宽表思想3.2 IntEnumField 与 JSONField 的存储语义3.3 text / keyword / date 三类字段选型3.4 异步客户端单例与依赖注入3.5 批量写入与幂等_id async_bulk3.6 N1 查询的批量化解法四、代码逻辑拆解4.1 ES 客户端单例与依赖桥接4.2 值序列化工具 _to_es_value4.3 宽表文档拼装 _build_job_document4.4 创建索引mapping 设计4.5 全量同步 insert-data-v24.6 职位写入接口 saveJobFastAPI Tortoise-ORM Elasticsearch 实战职位数据全量同步到 ES 的设计意识与踩坑复盘一、前言介绍1.1 功能定位本文只讲两件事的设计意识其一职位领域模型与写入接口怎么落地其二如何把 MySQL 里的职位数据批量、可靠地同步进 ES。1.2 数据模型总览三个实体之间的关系是同步逻辑的基础Job职位表 t_job └── enterprise_id : IntField ──┐ 手动整型关联不做外键级联 ↓ Enterprise企业主表 ── 1:1 ── EnterpriseInfo企业工商信息 └── industry : ForeignKeyField → Industry行业 Industry行业字典/ City城市字典要点Job与企业之间用的是enterprise_id普通整型字段而非ForeignKeyField这意味着同步时要手动按 ID 取企业且不会因为企业被删而级联掉职位。1.3 同步流程总览启动期lifespan 里初始化 ES 异步客户端单例 ↓ 建索引POST /es-data/create-index-v2 → 声明 mapping字段类型 IK 分词 ↓ 全量同步POST /es-data/insert-data-v2 Job.all() → 收集 enterprise_id批量取 Enterprise / EnterpriseInfo 进内存字典 → 每条职位拼成宽表文档 → async_bulk 批量写入_idjob.id 保证幂等二、环境准备2.1 依赖与 ES 客户端同步能力依赖官方elasticsearch异步客户端。连接信息从环境变量读取缺省回落到本机ES_HOSTos.getenv(ES_HOST,http://localhost:9200)es_client:AsyncElasticsearch|NoneNoneos.getenv(ES_HOST, ...)优先取环境变量便于不同环境切换 ES 地址es_client用模块级None占位后续做单例整个进程只建一条连接。2.2 索引与运行环境ES 必须安装IK 分词插件ik_max_word否则analyzer: ik_max_word建索引会报错本地开发可用单分片零副本number_of_shards: 1, number_of_replicas: 0职位索引取名boss_job_index_v2刻意与旧版boss_job_index分开便于对照学习。2.3 路由与生命周期挂载ES 客户端在应用启动期就初始化避免首次请求时再去建连asynccontextmanagerasyncdeflifespan(app:FastAPI):awaitTortoise.init(configTORTOISE_ORM,_enable_global_fallbackTrue)...awaitget_es_client()# 启动即建立 ES 连接单例yieldawaitTortoise.close_connections()lifespan是 FastAPI 的启动/关闭钩子在yield之前做的都是启动准备之后是优雅关闭这里只关了 TortoiseES 客户端关闭函数虽已备好但并未在此显式调用见问题排查 5.6。三、知识点讲解3.1 ORM 与 ES 的数据形态差异宽表思想MySQL 里职位、企业、工商信息、行业分表存储靠关联还原。ES 是文档型存储更适合把一次搜索要展示的所有字段拍平成一条文档避免搜索时再回查多表。这就是同步脚本把四张表拼成一份文档的根本动机。3.2 IntEnumField 与 JSONField 的存储语义statusfields.IntEnumField(enum_typeJobStatus,description0:草稿,1:招聘中...)department_idfields.IntEnumField(enum_typeDeptType,description所属部门)job_tagsfields.JSONField(defaultlist,description职位标签示例[五险一金,年终奖])IntEnumField数据库存的是整数但 ORM 层自动转成枚举对象写代码用JobStatus.RECRUITING比裸数字更安全JSONField数据库列里直接存 JSON 数组job_tags变成[五险一金,年终奖]进 ES 时映射成keyword多值字段。3.3 text / keyword / date 三类字段选型这是 mapping 设计的核心判断text ik_max_word职位名、公司名、行业、经营范围——需要中文全文检索keyword薪资、城市、标签、统一社会信用代码——用于精确匹配 / 聚合 / 筛选不分词integer / long状态、部门、企业 ID——用于数值范围与等值筛选date发布时间、认证时间——用于时间区间查询写入时用 ISO 字符串。3.4 异步客户端单例与依赖注入asyncdefget_es_client()-AsyncElasticsearch:globales_clientifes_clientisNone:es_clientAsyncElasticsearch(ES_HOST)logger.info(ES 客户端初始化成功)returnes_clientglobal es_client允许函数内修改模块级变量二次进入直接返回已建好的实例进程内复用一条连接避免每次请求都握手接口层通过Depends(es_client_depend)拿到它与鉴权依赖写法一致。3.5 批量写入与幂等_id async_bulkactions.append({_index:BOSS_JOB_INDEX_NAME_V2,_id:str(job.id),_source:document,})awaitasync_bulk(clientes_client,actionsactions,raise_on_errorFalse)_idstr(job.id)把职位主键设为文档 ID重复同步时覆盖旧文档而非新增天然幂等async_bulk一次性批量提交远比逐条index()快raise_on_errorFalse单条失败不中断整批返回值里能拿到错误列表做统计。3.6 N1 查询的批量化解法职位有 N 条若每条都现场查企业、查工商信息就是2N次查询N1 的变体。正确做法是先收集所有enterprise_id两批查完建字典内存里 O(1) 关联enterprise_idslist({job.enterprise_idforjobinjobsifjob.enterprise_idisnotNone})enterprisesawaitEnterprise.filter(id__inenterprise_ids).prefetch_related(city)enterprise_map{e.id:eforeinenterprises}集合去重避免重复查询prefetch_related(city)一次性把城市关联出来后面取enterprise.city.name不再发 SQLenterprise_map以企业 ID 为键的字典拼文档时直接enterprise_map.get(job.enterprise_id)。四、代码逻辑拆解4.1 ES 客户端单例与依赖桥接asyncdefes_client_depend()-AsyncElasticsearch:returnawaitget_es_client()这是 FastAPI 依赖函数作用是在路由签名里优雅注入 ES 客户端内部直接复用上面讲的单例保证全链路同一连接。4.2 值序列化工具 _to_es_valueES 只认 JSON 友好的类型ORM 里的datetime、date、枚举要先转换def_to_es_value(value:Any)-Any:ifvalueisNone:returnNoneifisinstance(value,datetime):returnvalue.isoformat()ifisinstance(value,date):returnvalue.isoformat()ifisinstance(value,Enum):returnvalue.valuereturnvalue第 1 行空值原样返回Nonemapping 里对应字段可空第 2–3 行datetime/date转 ISO 字符串ESdate字段才能识别第 4 行Enum转成.value一般是 int否则 ES 写入枚举对象会序列化失败最后一行其余类型str、int、list、dict、None原样透传。4.3 宽表文档拼装 _build_job_document这是同步的灵魂——把职位、企业、工商、行业压成一条扁平文档citygetattr(enterprise,city,None)ifenterpriseelseNoneindustrygetattr(enterprise_info,industry,None)ifenterprise_infoelseNonereturn{job_id:job.id,job_name:job.job_name,min_salary:job.min_salary,max_salary:job.max_salary,job_tags:job.job_tagsor[],status:_to_es_value(job.status),publish_time:_to_es_value(job.publish_time),enterprise_name:enterprise.enterprise_nameifenterpriseelseNone,enterprise_account_status:_to_es_value(enterprise.account_status)ifenterpriseelseNone,enterpriseInfo_company_scale:_to_es_value(enterprise_info.company_scale)ifenterprise_infoelseNone,industry_id:industry.idifindustryelseNone,industry_name:industry.nameifindustryelseNone,}getattr(enterprise, city, None)防御式取值企业对象没有city属性时不抛异常job_tags or []标签为空时给空数组避免 ES 收到None与keyword多值类型冲突if enterprise else None企业缺失时整组企业字段填None保证单条脏数据不会打断整批字段名带enterpriseInfo_前缀是为了和 mapping 一一对应也能直观区分属于企业详情凡是status、publish_time、company_scale等枚举/时间字段一律过_to_es_value。4.4 创建索引mapping 设计建索引时声明每个字段类型与中文分词器这些是经旧版踩坑后补全的job_name:{type:text,analyzer:ik_max_word},enterprise_name:{type:text,analyzer:ik_max_word},work_location:{type:keyword},min_salary:{type:keyword},status:{type:integer},publish_time:{type:date},enterpriseInfo_business_scope:{type:text,analyzer:ik_max_word},职位名、公司名、经营范围用text ik_max_word支持中文全文检索城市、薪资用keyword因为要做精确筛选与聚合不能分词status用integer和IntEnumField存的整数对齐enterpriseInfo_business_scope旧版 mapping 漏声明却仍在写入导致该字段无法被检索这里显式补回text类型建索引前先indices.exists判断已存在则直接返回避免重复创建报错。4.5 全量同步 insert-data-v2核心流程先校验索引存在再批量取数再组装 bulk最后一次性写入。ifnotawaites_client.indices.exists(indexBOSS_JOB_INDEX_NAME_V2):return{code:0,message:f索引{BOSS_JOB_INDEX_NAME_V2}不存在请先调用 /es-data/create-index-v2}jobsawaitJob.all()ifnotjobs:return{code:1,message:没有可同步的职位,data:{success:0,skip:0}}第一步先确认目标索引已建否则 bulk 时才报错浪费前面查库的开销Job.all()一次取出全部职位空集合提前返回避免后面空跑。forjobinjobs:enterpriseenterprise_map.get(job.enterprise_id)enterprise_infoenterprise_info_map.get(job.enterprise_id)ifenterpriseisNone:skip_count1logger.warning(f同步跳过职位 id{job.id}关联企业 id{job.enterprise_id}不存在)continueifenterprise_infoisNone:logger.warning(f职位 id{job.id}无企业详情将写入空的 enterpriseInfo_* 字段)document_build_job_document(job,enterprise,enterprise_info)actions.append({_index:BOSS_JOB_INDEX_NAME_V2,_id:str(job.id),_source:document})enterprise_map.get(...)内存字典取值O(1)无额外 SQL企业主数据缺失则continue跳过宽表缺核心信息写入也搜不到意义不大并记日志企业详情缺失允许继续只是enterpriseInfo_*字段为None每条文档带_idstr(job.id)这是幂等覆盖的关键全部 action 先攒进列表最后统一async_bulk。success_count,errorsawaitasync_bulk(clientes_client,actionsactions,raise_on_errorFalse,)error_countlen(errors)ifisinstance(errors,list)else0return{code:1,message:数据同步完成,data:{index:BOSS_JOB_INDEX_NAME_V2,job_total:len(jobs),success:success_count,skip:skip_count,error:error_count},}success_count是成功条数errors是失败详情列表raise_on_errorFalse让部分失败不影响整体返回里把成功/跳过/失败三类数量都带上便于对账。4.6 职位写入接口 saveJob职位落库时企业 ID 不是前端传的而是从登录的招聘团队反查出来避免越权挂到别家企业名下staticmethodasyncdefsaveJob(job:JobCreateRequest,time_id:int):timeawaitRecruitTeam.filter(idtime_id).first()awaitJob.create(**job.dict(),enterprise_idtime.enterprise_id,recruit_team_idtime_id,publish_timenow(),)time_id来自 JWT 鉴权依赖get_job_info代表当前招聘团队RecruitTeam.filter(idtime_id).first()查出团队取其enterprise_id**job.dict()把请求体字段展开成建表参数省去逐字段赋值enterprise_idtime.enterprise_id企业归属由服务端决定前端无法伪造publish_timenow()写入发布时间后续同步进 ES 的date字段。