戰(zhàn)指南:3 步讓任務(wù)不再重復(fù)計(jì)算)
Prefect 緩存策略實(shí)戰(zhàn)指南3 步讓任務(wù)不再重復(fù)計(jì)算【免費(fèi)下載鏈接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.項(xiàng)目地址: https://gitcode.com/GitHub_Trending/pr/prefectPrefect 是一個用 Python 構(gòu)建數(shù)據(jù)管道的任務(wù)編排框架它的任務(wù)緩存Cache Policy能幫你把算過的任務(wù)結(jié)果存下來、下次直接取從而省掉大量重復(fù)計(jì)算。這篇文章圍繞三個問題展開你的任務(wù)該不該緩存、Prefect 緩存策略怎么配、配置之后怎么驗(yàn)證命中、出了問題怎么排。全文以可操作步驟為主跟著做就能跑通。一、開場同一條數(shù)據(jù)一天被算了三遍想象這樣一個畫面你的夜間 ETL 流程每小時跑一次其中拉取昨日訂單匯總這個任務(wù)每次都要連數(shù)據(jù)庫、跑兩分鐘聚合。凌晨 2 點(diǎn)跑一次、4 點(diǎn)跑一次、6 點(diǎn)再跑一次——輸入?yún)?shù)一模一樣結(jié)果一模一樣機(jī)器空轉(zhuǎn)六分鐘。你會遇到的情況就是這類任務(wù)重復(fù)執(zhí)行但產(chǎn)出沒有任何變化。在 Prefect 里task默認(rèn)每次都真跑。如果你給它配上緩存策略第二次執(zhí)行時框架會直接拿出上次的結(jié)果任務(wù)耗時從兩分鐘變成幾秒。判斷標(biāo)準(zhǔn)很簡單這個任務(wù)換個時間點(diǎn)再算一遍結(jié)果會不會變。如果答案是不會它就有緩存價值。那么哪些任務(wù)符合這個條件下一章給一張自檢表。二、該不該緩存先做判斷再動手緩存不是越用越好。給一個有副作用或結(jié)果隨時變化的任務(wù)貼緩存反而會讓你讀到過期數(shù)據(jù)。動手之前先按下面這張表自檢你的任務(wù)類型任務(wù)特征適合緩存典型例子冪等同樣的輸入永遠(yuǎn)得到同樣的輸出適合讀數(shù)、聚合、格式轉(zhuǎn)換輸入穩(wěn)定參數(shù)在短期內(nèi)不會變適合固定表名、固定查詢條件計(jì)算成本高或調(diào)用外部 API 有配額適合大表 join、付費(fèi)接口查詢結(jié)果有時效性分鐘級就變謹(jǐn)慎配短過期時間實(shí)時價格、庫存有副作用寫庫、發(fā)消息、扣款不適合任何執(zhí)行一次就該只執(zhí)行一次的操作隨機(jī)性內(nèi)部用了隨機(jī)數(shù)、當(dāng)前時間戳不適合除非把時間固定進(jìn)輸入抽樣、打時間戳的報(bào)表怎么用這張表挑出適合一欄里的任務(wù)去配緩存不適合的一欄保持默認(rèn)每次執(zhí)行謹(jǐn)慎的一欄配緩存但把過期時間壓短。驗(yàn)證這一步是否做對也很直接給一個冪等任務(wù)開緩存、給一個寫庫任務(wù)不開跑兩遍流程前者第二次秒回后者兩次都真實(shí)執(zhí)行——符合預(yù)期說明判斷對了。三、看懂機(jī)制一次 Prefect 緩存的一生Prefect 緩存的源碼實(shí)現(xiàn)分布在兩處任務(wù)側(cè)的task參數(shù)src/prefect/tasks.py和編排側(cè)的 核心策略源碼。在CoreTaskPolicy的規(guī)則優(yōu)先級里CacheRetrieval排在最前、CacheInsertion排在后段——翻譯過來就是先查緩存再寫緩存。你可以把緩存理解成給算過的結(jié)果貼標(biāo)簽標(biāo)簽緩存鍵一樣就直接取下上次貼的標(biāo)簽里的東西。按時間線看一次任務(wù)的完整生命周期未命中任務(wù)啟動前CacheRetrieval規(guī)則拿緩存鍵去查庫。第一次跑庫里沒有對應(yīng)記錄放行執(zhí)行。執(zhí)行任務(wù)真實(shí)運(yùn)行跑完進(jìn)入成功終態(tài)。寫入CacheInsertion規(guī)則把緩存鍵 → 這次的狀態(tài)/結(jié)果寫進(jìn)數(shù)據(jù)庫。存儲表是task_run_state_cache定義在 ORM 模型 里每條記錄包含緩存鍵、關(guān)聯(lián)的狀態(tài) ID 和創(chuàng)建時間。命中下次同鍵任務(wù)啟動第 1 步查到記錄任務(wù)直接跳到成功態(tài)不再執(zhí)行函數(shù)體。過期/失效如果配置了過期時間寫入時超過時長的舊記錄會被忽略或者你改了輸入、改了版本號鍵變了舊緩存自然作廢。驗(yàn)證機(jī)制生效的辦法跑同一流程兩遍看第二個任務(wù)運(yùn)行是否沒有實(shí)際執(zhí)行函數(shù)日志里沒有函數(shù)內(nèi)的打印狀態(tài)卻直接是 Completed并且任務(wù)記錄上能看到cache_key。四、動手三步配出可用的任務(wù)緩存第一步給任務(wù)生成緩存鍵沒有緩存鍵就沒有緩存。最小可用配置是套用內(nèi)置的task_input_hash它會對任務(wù)名、函數(shù)代碼和全部入?yún)⒆龉!囊粋€參數(shù)鍵就變緩存自動失效from prefect import task from prefect.tasks import task_input_hash task(cache_key_fntask_input_hash) def build_summary(date: str, region: str): # 你的查詢邏輯 return run_query(...)這一行cache_key_fntask_input_hash就是全部開關(guān)。跑兩遍同參數(shù)流程驗(yàn)證第二遍該任務(wù)應(yīng)秒級結(jié)束。第二步設(shè)置緩存多久過期結(jié)果有時效的任務(wù)加一個cache_expiration即可。經(jīng)驗(yàn)值是不超過你業(yè)務(wù)上能容忍的數(shù)據(jù)延遲from datetime import timedelta task( cache_key_fntask_input_hash, cache_expirationtimedelta(hours6), ) def fetch_price(): ...6 小時內(nèi)的重復(fù)請求走緩存超過 6 小時自動重新計(jì)算。想驗(yàn)證等過期后或臨時改成幾分鐘再跑確認(rèn)任務(wù)重新真實(shí)執(zhí)行。第三步規(guī)劃好失效手段有三種讓舊緩存作廢的方式按常用程度排改輸入task_input_hash已覆蓋參數(shù)一變自動失效換版本改邏輯但輸入沒變時用task(version2.0)版本號參與哈希升級即全部失效清空重來極端情況直接刪庫里的task_run_state_cache記錄謹(jǐn)慎使用。三步走完你的任務(wù)就有了自動命中、按時間過期、按版本換代的完整閉環(huán)。五、排障命中率和鍵沖突的現(xiàn)場排查緩存配好之后最常見的兩類問題是該命中沒命中和命中了錯誤的數(shù)據(jù)。按下面的流程走不用靠猜癥狀 A命中率低任務(wù)總在重跑定位先看是不是參數(shù)看起來一樣、實(shí)際不一樣——task_input_hash會對入?yún)⒆鰢?yán)格哈希字典順序外的差異、datetime和字符串的差異都會改變鍵打印兩次運(yùn)行的cache_key對比一下。定位確認(rèn)函數(shù)體有沒有被改過——task_input_hash會把函數(shù)字節(jié)碼算進(jìn)鍵改了一行代碼等于全量失效這是特性不是故障。處理把每次都變的參數(shù)如當(dāng)前時間從入?yún)⑴策M(jìn)函數(shù)內(nèi)部或改用自定義cache_key_fn只哈希關(guān)鍵參數(shù)。癥狀 B命中了但數(shù)據(jù)是錯的鍵沖突定位不同任務(wù)或不同環(huán)境的同名函數(shù)撞了緩存鍵。task_input_hash含任務(wù)名通常不會跨任務(wù)沖突自己寫的簡單cache_key_fn最容易漏項(xiàng)。處理給鍵加前綴或命名空間例如把任務(wù)名、版本號、環(huán)境名都拼進(jìn)去做到一個鍵只屬于一個任務(wù)的一個版本。驗(yàn)證修復(fù)后清掉舊鍵對應(yīng)的記錄重跑確認(rèn)新舊環(huán)境各走各的緩存。癥狀 C緩存越攢越大定位查task_run_state_cache表記錄數(shù)是否持續(xù)上漲多半是沒設(shè)過期時間。處理給任務(wù)補(bǔ)上cache_expiration或定期清理無引用的舊記錄。排障的總原則先比鍵再比版本最后才懷疑框架。六、進(jìn)階按環(huán)境和條件切換緩存行為你經(jīng)常需要開發(fā)環(huán)境每次都真跑、生產(chǎn)環(huán)境盡量走緩存這類行為差異。做法是自己寫一個cache_key_fn讓它根據(jù)條件返回鍵或返回None返回None表示這次不緩存import os from prefect.tasks import task_input_hash def env_cache_key(context, arguments): if os.environ.get(PREFECT_ENV) dev: return None return task_input_hash(context, arguments) task(cache_key_fnenv_cache_key) def etl_step(...): ...注意簽名是(context, arguments)context里能拿到任務(wù)運(yùn)行上下文arguments是入?yún)⒆值洹M瑯拥膶懛ㄒ部梢宰龀砂醋鈶?、按?shù)據(jù)分區(qū)、按開關(guān)動態(tài)切換。驗(yàn)證方式在 dev 環(huán)境跑兩遍確認(rèn)每次都真執(zhí)行切到生產(chǎn)環(huán)境再跑兩遍第二遍應(yīng)命中。這樣同一份代碼行為隨環(huán)境自動切換不需要維護(hù)兩套任務(wù)定義。七、收尾上線前自查清單 延伸資源把緩存策略用到生產(chǎn)之前逐項(xiàng)勾一遍只給冪等、輸入穩(wěn)定的任務(wù)開了緩存副作用任務(wù)保持每跑必執(zhí)行緩存鍵用task_input_hash或包含任務(wù)名 版本 關(guān)鍵入?yún)⒌淖远x鍵無跨任務(wù)沖突有時效的結(jié)果都設(shè)了cache_expiration時長不超過業(yè)務(wù)可容忍延遲改邏輯時用version升級觸發(fā)全量失效而不是靠祈禱開發(fā)環(huán)境禁緩存、生產(chǎn)啟用緩存的行為已驗(yàn)證跑了兩遍流程確認(rèn)第二遍命中緩存且結(jié)果正確知道去哪里查task_run_state_cache表來排查鍵沖突延伸資源均為倉庫內(nèi)文件可相對路徑直接打開緩存鍵生成與任務(wù)參數(shù)定義src/prefect/tasks.py檢索/寫入規(guī)則的編排順序src/prefect/server/orchestration/core_policy.py存儲表結(jié)構(gòu)src/prefect/server/database/orm_models.py官方文檔索引docs/緩存這件事Prefect 把何時命中、何時寫入、何時失效都收斂到了任務(wù)參數(shù)和編排規(guī)則里。你只需要回答三個問題鍵怎么生成、過期多久、什么時候換版本。答對了重復(fù)計(jì)算就消失了?!久赓M(fèi)下載鏈接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.項(xiàng)目地址: https://gitcode.com/GitHub_Trending/pr/prefect創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考