:職位數(shù)據(jù)全量同步到 ES 的設(shè)計)
項目實踐FastAPI Elasticsearch 職位數(shù)據(jù)全量同步FastAPI Tortoise-ORM Elasticsearch 實戰(zhàn)職位數(shù)據(jù)全量同步到 ES 的設(shè)計意識與踩坑復(fù)盤一、前言介紹1.1 功能定位1.2 數(shù)據(jù)模型總覽1.3 同步流程總覽二、環(huán)境準(zhǔn)備2.1 依賴與 ES 客戶端2.2 索引與運行環(huán)境2.3 路由與生命周期掛載三、知識點講解3.1 ORM 與 ES 的數(shù)據(jù)形態(tài)差異寬表思想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 創(chuàng)建索引mapping 設(shè)計4.5 全量同步 insert-data-v24.6 職位寫入接口 saveJobFastAPI Tortoise-ORM Elasticsearch 實戰(zhàn)職位數(shù)據(jù)全量同步到 ES 的設(shè)計意識與踩坑復(fù)盤一、前言介紹1.1 功能定位本文只講兩件事的設(shè)計意識其一職位領(lǐng)域模型與寫入接口怎么落地其二如何把 MySQL 里的職位數(shù)據(jù)批量、可靠地同步進(jìn) ES。1.2 數(shù)據(jù)模型總覽三個實體之間的關(guān)系是同步邏輯的基礎(chǔ)Job職位表 t_job └── enterprise_id : IntField ──┐ 手動整型關(guān)聯(lián)不做外鍵級聯(lián) ↓ Enterprise企業(yè)主表 ── 1:1 ── EnterpriseInfo企業(yè)工商信息 └── industry : ForeignKeyField → Industry行業(yè) Industry行業(yè)字典/ City城市字典要點Job與企業(yè)之間用的是enterprise_id普通整型字段而非ForeignKeyField這意味著同步時要手動按 ID 取企業(yè)且不會因為企業(yè)被刪而級聯(lián)掉職位。1.3 同步流程總覽啟動期lifespan 里初始化 ES 異步客戶端單例 ↓ 建索引POST /es-data/create-index-v2 → 聲明 mapping字段類型 IK 分詞 ↓ 全量同步POST /es-data/insert-data-v2 Job.all() → 收集 enterprise_id批量取 Enterprise / EnterpriseInfo 進(jìn)內(nèi)存字典 → 每條職位拼成寬表文檔 → async_bulk 批量寫入_idjob.id 保證冪等二、環(huán)境準(zhǔn)備2.1 依賴與 ES 客戶端同步能力依賴官方elasticsearch異步客戶端。連接信息從環(huán)境變量讀取缺省回落到本機(jī)ES_HOSTos.getenv(ES_HOST,http://localhost:9200)es_client:AsyncElasticsearch|NoneNoneos.getenv(ES_HOST, ...)優(yōu)先取環(huán)境變量便于不同環(huán)境切換 ES 地址es_client用模塊級None占位后續(xù)做單例整個進(jìn)程只建一條連接。2.2 索引與運行環(huán)境ES 必須安裝IK 分詞插件ik_max_word否則analyzer: ik_max_word建索引會報錯本地開發(fā)可用單分片零副本number_of_shards: 1, number_of_replicas: 0職位索引取名boss_job_index_v2刻意與舊版boss_job_index分開便于對照學(xué)習(xí)。2.3 路由與生命周期掛載ES 客戶端在應(yīng)用啟動期就初始化避免首次請求時再去建連asynccontextmanagerasyncdeflifespan(app:FastAPI):awaitTortoise.init(configTORTOISE_ORM,_enable_global_fallbackTrue)...awaitget_es_client()# 啟動即建立 ES 連接單例yieldawaitTortoise.close_connections()lifespan是 FastAPI 的啟動/關(guān)閉鉤子在yield之前做的都是啟動準(zhǔn)備之后是優(yōu)雅關(guān)閉這里只關(guān)了 TortoiseES 客戶端關(guān)閉函數(shù)雖已備好但并未在此顯式調(diào)用見問題排查 5.6。三、知識點講解3.1 ORM 與 ES 的數(shù)據(jù)形態(tài)差異寬表思想MySQL 里職位、企業(yè)、工商信息、行業(yè)分表存儲靠關(guān)聯(lián)還原。ES 是文檔型存儲更適合把一次搜索要展示的所有字段拍平成一條文檔避免搜索時再回查多表。這就是同步腳本把四張表拼成一份文檔的根本動機(jī)。3.2 IntEnumField 與 JSONField 的存儲語義statusfields.IntEnumField(enum_typeJobStatus,description0:草稿,1:招聘中...)department_idfields.IntEnumField(enum_typeDeptType,description所屬部門)job_tagsfields.JSONField(defaultlist,description職位標(biāo)簽示例[五險一金,年終獎])IntEnumField數(shù)據(jù)庫存的是整數(shù)但 ORM 層自動轉(zhuǎn)成枚舉對象寫代碼用JobStatus.RECRUITING比裸數(shù)字更安全JSONField數(shù)據(jù)庫列里直接存 JSON 數(shù)組job_tags變成[五險一金,年終獎]進(jìn) ES 時映射成keyword多值字段。3.3 text / keyword / date 三類字段選型這是 mapping 設(shè)計的核心判斷text ik_max_word職位名、公司名、行業(yè)、經(jīng)營范圍——需要中文全文檢索keyword薪資、城市、標(biāo)簽、統(tǒng)一社會信用代碼——用于精確匹配 / 聚合 / 篩選不分詞integer / long狀態(tài)、部門、企業(yè) ID——用于數(shù)值范圍與等值篩選date發(fā)布時間、認(rèn)證時間——用于時間區(qū)間查詢寫入時用 ISO 字符串。3.4 異步客戶端單例與依賴注入asyncdefget_es_client()-AsyncElasticsearch:globales_clientifes_clientisNone:es_clientAsyncElasticsearch(ES_HOST)logger.info(ES 客戶端初始化成功)returnes_clientglobal es_client允許函數(shù)內(nèi)修改模塊級變量二次進(jìn)入直接返回已建好的實例進(jìn)程內(nèi)復(fù)用一條連接避免每次請求都握手接口層通過Depends(es_client_depend)拿到它與鑒權(quán)依賴寫法一致。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)把職位主鍵設(shè)為文檔 ID重復(fù)同步時覆蓋舊文檔而非新增天然冪等async_bulk一次性批量提交遠(yuǎn)比逐條index()快raise_on_errorFalse單條失敗不中斷整批返回值里能拿到錯誤列表做統(tǒng)計。3.6 N1 查詢的批量化解法職位有 N 條若每條都現(xiàn)場查企業(yè)、查工商信息就是2N次查詢N1 的變體。正確做法是先收集所有enterprise_id兩批查完建字典內(nèi)存里 O(1) 關(guān)聯(lián)enterprise_idslist({job.enterprise_idforjobinjobsifjob.enterprise_idisnotNone})enterprisesawaitEnterprise.filter(id__inenterprise_ids).prefetch_related(city)enterprise_map{e.id:eforeinenterprises}集合去重避免重復(fù)查詢prefetch_related(city)一次性把城市關(guān)聯(lián)出來后面取enterprise.city.name不再發(fā) SQLenterprise_map以企業(yè) ID 為鍵的字典拼文檔時直接enterprise_map.get(job.enterprise_id)。四、代碼邏輯拆解4.1 ES 客戶端單例與依賴橋接asyncdefes_client_depend()-AsyncElasticsearch:returnawaitget_es_client()這是 FastAPI 依賴函數(shù)作用是在路由簽名里優(yōu)雅注入 ES 客戶端內(nèi)部直接復(fù)用上面講的單例保證全鏈路同一連接。4.2 值序列化工具 _to_es_valueES 只認(rèn) JSON 友好的類型ORM 里的datetime、date、枚舉要先轉(zhuǎn)換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 里對應(yīng)字段可空第 2–3 行datetime/date轉(zhuǎn) ISO 字符串ESdate字段才能識別第 4 行Enum轉(zhuǎn)成.value一般是 int否則 ES 寫入枚舉對象會序列化失敗最后一行其余類型str、int、list、dict、None原樣透傳。4.3 寬表文檔拼裝 _build_job_document這是同步的靈魂——把職位、企業(yè)、工商、行業(yè)壓成一條扁平文檔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)防御式取值企業(yè)對象沒有city屬性時不拋異常job_tags or []標(biāo)簽為空時給空數(shù)組避免 ES 收到None與keyword多值類型沖突if enterprise else None企業(yè)缺失時整組企業(yè)字段填None保證單條臟數(shù)據(jù)不會打斷整批字段名帶enterpriseInfo_前綴是為了和 mapping 一一對應(yīng)也能直觀區(qū)分屬于企業(yè)詳情凡是status、publish_time、company_scale等枚舉/時間字段一律過_to_es_value。4.4 創(chuàng)建索引mapping 設(shè)計建索引時聲明每個字段類型與中文分詞器這些是經(jīng)舊版踩坑后補全的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},職位名、公司名、經(jīng)營范圍用text ik_max_word支持中文全文檢索城市、薪資用keyword因為要做精確篩選與聚合不能分詞status用integer和IntEnumField存的整數(shù)對齊enterpriseInfo_business_scope舊版 mapping 漏聲明卻仍在寫入導(dǎo)致該字段無法被檢索這里顯式補回text類型建索引前先indices.exists判斷已存在則直接返回避免重復(fù)創(chuàng)建報錯。4.5 全量同步 insert-data-v2核心流程先校驗索引存在再批量取數(shù)再組裝 bulk最后一次性寫入。ifnotawaites_client.indices.exists(indexBOSS_JOB_INDEX_NAME_V2):return{code:0,message:f索引{BOSS_JOB_INDEX_NAME_V2}不存在請先調(diào)用 /es-data/create-index-v2}jobsawaitJob.all()ifnotjobs:return{code:1,message:沒有可同步的職位,data:{success:0,skip:0}}第一步先確認(rèn)目標(biāo)索引已建否則 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}關(guān)聯(lián)企業(yè) id{job.enterprise_id}不存在)continueifenterprise_infoisNone:logger.warning(f職位 id{job.id}無企業(yè)詳情將寫入空的 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(...)內(nèi)存字典取值O(1)無額外 SQL企業(yè)主數(shù)據(jù)缺失則continue跳過寬表缺核心信息寫入也搜不到意義不大并記日志企業(yè)詳情缺失允許繼續(xù)只是enterpriseInfo_*字段為None每條文檔帶_idstr(job.id)這是冪等覆蓋的關(guān)鍵全部 action 先攢進(jìn)列表最后統(tǒng)一async_bulk。success_count,errorsawaitasync_bulk(clientes_client,actionsactions,raise_on_errorFalse,)error_countlen(errors)ifisinstance(errors,list)else0return{code:1,message:數(shù)據(jù)同步完成,data:{index:BOSS_JOB_INDEX_NAME_V2,job_total:len(jobs),success:success_count,skip:skip_count,error:error_count},}success_count是成功條數(shù)errors是失敗詳情列表raise_on_errorFalse讓部分失敗不影響整體返回里把成功/跳過/失敗三類數(shù)量都帶上便于對賬。4.6 職位寫入接口 saveJob職位落庫時企業(yè) ID 不是前端傳的而是從登錄的招聘團(tuán)隊反查出來避免越權(quán)掛到別家企業(yè)名下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 鑒權(quán)依賴get_job_info代表當(dāng)前招聘團(tuán)隊RecruitTeam.filter(idtime_id).first()查出團(tuán)隊取其enterprise_id**job.dict()把請求體字段展開成建表參數(shù)省去逐字段賦值enterprise_idtime.enterprise_id企業(yè)歸屬由服務(wù)端決定前端無法偽造publish_timenow()寫入發(fā)布時間后續(xù)同步進(jìn) ES 的date字段。