詳解:限制任務(wù)執(zhí)行并行度與多槽位調(diào)度實(shí)戰(zhàn))
Apache Airflow 中的 Pools資源池詳解限制任務(wù)執(zhí)行并行度與多槽位調(diào)度實(shí)戰(zhàn)【免費(fèi)下載鏈接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows項(xiàng)目地址: https://gitcode.com/GitHub_Trending/ai/airflow導(dǎo)讀當(dāng) DAG 中的多個(gè)任務(wù)在相近時(shí)間集中觸發(fā)時(shí)數(shù)據(jù)庫(kù)、外部 API 或遺留系統(tǒng)很容易因瞬時(shí)并發(fā)過(guò)高而過(guò)載甚至被打垮。Airflow 的 Pools資源池正是用于限制任意任務(wù)集合的執(zhí)行并行度的機(jī)制——通過(guò)給池命名并分配固定數(shù)量的 worker 槽位slots將同時(shí)運(yùn)行的任務(wù)數(shù)牢牢控制在目標(biāo)系統(tǒng)可承受的范圍內(nèi)。本文以 Apache Airflow 倉(cāng)庫(kù)中的官方管理文檔為基礎(chǔ)完整講解池的創(chuàng)建與管理、通過(guò)pool參數(shù)綁定任務(wù)、利用pool_slots讓任務(wù)按“計(jì)算權(quán)重”占用多個(gè)槽位并結(jié)合調(diào)度器與數(shù)據(jù)模型源碼說(shuō)明槽位統(tǒng)計(jì)、排隊(duì)與放行的底層原理。讀完你將能獨(dú)立設(shè)計(jì)一套基于資源池的限流調(diào)度方案。Pools 能解決什么問(wèn)題在 Airflow 中DAG 的調(diào)度決定了任務(wù)“何時(shí)應(yīng)該運(yùn)行”但默認(rèn)情況下并不天然限制“同一時(shí)刻有多少個(gè)任務(wù)并行壓向同一個(gè)下游系統(tǒng)”。當(dāng)以下場(chǎng)景疊加時(shí)問(wèn)題會(huì)迅速放大多個(gè) DAG 共享同一個(gè)數(shù)據(jù)庫(kù)實(shí)例或同一組 API 額度大量任務(wù)在同一個(gè)時(shí)間窗如整點(diǎn)集中變得可運(yùn)行不同任務(wù)的計(jì)算負(fù)載差異懸殊個(gè)別任務(wù)會(huì)瞬間吃滿外部系統(tǒng)資源。Airflow Pools 提供的是一個(gè)應(yīng)用層面的并發(fā)閘門(mén)它允許把一批“會(huì)命中同一目標(biāo)系統(tǒng)”的任務(wù)歸入同一個(gè)池并限定該池最多同時(shí)占用多少槽位。任務(wù)的常規(guī)調(diào)度照常進(jìn)行狀態(tài)流轉(zhuǎn)、依賴關(guān)系、重試等都不受影響但一旦池的容量被占滿可運(yùn)行的任務(wù)會(huì)進(jìn)入排隊(duì)狀態(tài)在 UI 上顯示為queued當(dāng)槽位釋放后再依據(jù)任務(wù)的 Priority Weight優(yōu)先級(jí)權(quán)重 及其后代任務(wù)的權(quán)重依次放行。這一點(diǎn)與“限制并發(fā)但保留有序執(zhí)行”的運(yùn)維訴求精確對(duì)應(yīng)。在 UI 中創(chuàng)建與管理 Pools池的列表在 Web UI 中管理入口為Menu - Admin - Pools。在這里可以為每個(gè)池指定Pool名稱后續(xù)在 DAG 代碼中通過(guò)pool參數(shù)引用Slots槽位數(shù)該池允許同時(shí)被占用的最大槽位數(shù)量決定并發(fā)上限D(zhuǎn)escription描述便于運(yùn)維人員識(shí)別該池的用途如面向哪個(gè)下游系統(tǒng)是否把延遲任務(wù)計(jì)入槽位占用創(chuàng)建時(shí)可勾選是否在計(jì)算被占用槽位時(shí)把 延遲deferred任務(wù) 一并納入。其中“把 deferred 任務(wù)計(jì)入占用”是一個(gè)容易被忽略但重要的開(kāi)關(guān)。使用defer機(jī)制把任務(wù)轉(zhuǎn)成 triggerer 管理的延遲態(tài)后任務(wù)實(shí)例本身并不占用 worker但它仍占著池里的一個(gè)位置還是已經(jīng)釋放、把槽位讓給其它排隊(duì)任務(wù)這個(gè)開(kāi)關(guān)正是用來(lái)控制這一行為的。從數(shù)據(jù)模型看池在元數(shù)據(jù)庫(kù)中的表名為slot_pool見(jiàn) pool.py每個(gè)池記錄包含pool名稱唯一約束最長(zhǎng) 256 字符slots槽位總數(shù)代碼中以-1表示無(wú)限槽位description自由文本描述include_deferred布爾值決定DEFERRED狀態(tài)任務(wù)是否占用槽位team_name可選的所屬團(tuán)隊(duì)通過(guò)外鍵關(guān)聯(lián)到team.name團(tuán)隊(duì)刪除時(shí)置空。因此“給池分配名字和槽位”本質(zhì)上是對(duì)這張表做一次插入或更新而調(diào)度器在每一輪調(diào)度時(shí)都會(huì)實(shí)時(shí)讀取這些配置來(lái)決定放行哪些任務(wù)。將任務(wù)關(guān)聯(lián)到某個(gè) Pool創(chuàng)建好池之后在任務(wù)中使用pool參數(shù)即可完成關(guān)聯(lián)。官方文檔給出的典型示例是把批處理任務(wù)放進(jìn)一個(gè)面向消息聚合管道的池aggregate_db_message_job BashOperator( task_idaggregate_db_message_job, execution_timeouttimedelta(hours3), poolep_data_pipeline_db_msg_agg, bash_commandaggregate_db_message_job_cmd, dagdag, ) aggregate_db_message_job.set_upstream(wait_for_empty_queue)把a(bǔ)ggregate_db_message_job放入名為ep_data_pipeline_db_msg_agg的池后調(diào)度器只會(huì)從這個(gè)池的剩余可用槽位中“放行”該任務(wù)。槽位被占滿時(shí)所有可運(yùn)行但拿不到槽位的任務(wù)進(jìn)入queued狀態(tài)并持續(xù)等待隨著運(yùn)行中任務(wù)結(jié)束、槽位釋放排隊(duì)的任務(wù)會(huì)按優(yōu)先級(jí)權(quán)重依次被調(diào)度執(zhí)行權(quán)重的計(jì)算規(guī)則可參考 Priority Weight 文檔。值得注意的行為細(xì)節(jié)是任務(wù)在池耗盡時(shí)并不會(huì)失敗或超時(shí)它只是被排隊(duì)。所以一個(gè)池的容量設(shè)計(jì)應(yīng)當(dāng)結(jié)合任務(wù)平均執(zhí)行時(shí)長(zhǎng)與調(diào)度頻率綜合評(píng)估——若池太小而排入任務(wù)過(guò)多任務(wù)會(huì)長(zhǎng)時(shí)間停留在 queued 狀態(tài)進(jìn)而拉長(zhǎng)整個(gè)數(shù)據(jù)管道的端到端延遲。默認(rèn)池 default_pool如果沒(méi)有給任務(wù)顯式指定池任務(wù)會(huì)被自動(dòng)分配到名為default_pool的默認(rèn)池。默認(rèn)池初始化時(shí)有128 個(gè)槽位可以通過(guò) UI 或 CLI 修改槽位數(shù)量但不能被刪除。在源碼層面默認(rèn)池名稱以常量DEFAULT_POOL_NAME default_pool的形式定義在 pool.py 中并通過(guò)數(shù)據(jù)庫(kù)遷移在初始化時(shí)寫(xiě)入元數(shù)據(jù)庫(kù)。刪除池的邏輯對(duì)默認(rèn)池做了硬性保護(hù)staticmethod provide_session def delete_pool(name: str, *, session: Session NEW_SESSION) - Pool: Delete pool by a given name. if name Pool.DEFAULT_POOL_NAME: raise AirflowException(f{Pool.DEFAULT_POOL_NAME} cannot be deleted) ...也就是說(shuō)即使airflow pools delete default_pool也會(huì)被拋出的AirflowException拒絕。生產(chǎn)實(shí)踐中常見(jiàn)的做法是把default_pool的 128 個(gè)槽位調(diào)小甚至調(diào)到很小以“逼著”開(kāi)發(fā)者給每個(gè)會(huì)沖擊外部系統(tǒng)的任務(wù)顯式指定業(yè)務(wù)池避免海量未指定池的任務(wù)默認(rèn)并行打滿 128 路。用 pool_slots 讓任務(wù)占用多個(gè)槽位默認(rèn)情況下每個(gè)任務(wù)實(shí)例占用 1 個(gè)池槽位。但對(duì)于計(jì)算負(fù)載差異很大的任務(wù)組一律按 1 個(gè)槽位計(jì)并不公平。Airflow 提供了pool_slots參數(shù)允許單個(gè)任務(wù)在運(yùn)行時(shí)占用多個(gè)槽位。官方文檔用一個(gè)maintenance池共 2 個(gè)槽位說(shuō)明了它的價(jià)值BashOperator( task_idheavy_task, bash_commandbash backup_data.sh, pool_slots2, poolmaintenance, ) BashOperator( task_idlight_task1, bash_commandbash check_files.sh, pool_slots1, poolmaintenance, ) BashOperator( task_idlight_task2, bash_commandbash remove_files.sh, pool_slots1, poolmaintenance, )在這個(gè)例子中heavy_task配置占用 2 個(gè)槽位因此只要它處于運(yùn)行狀態(tài)就會(huì)耗盡maintenance池的全部 2 個(gè)槽位兩個(gè) light 任務(wù)必須排隊(duì)等待它結(jié)束反過(guò)來(lái)light_task1與light_task2各自只占 1 個(gè)槽位可以并發(fā)運(yùn)行而heavy_task需要等到兩個(gè)槽位同時(shí)空閑才會(huì)啟動(dòng)。這里的等價(jià)關(guān)系是在資源占用意義上一個(gè)占用 2 個(gè)槽位的 heavy 任務(wù) ≈ 兩個(gè)并發(fā)運(yùn)行的 light 任務(wù)。這種“按權(quán)重計(jì)槽”的設(shè)計(jì)直接防止了“一個(gè)重任務(wù)與一個(gè)輕任務(wù)并發(fā)運(yùn)行”時(shí)把系統(tǒng)資源瞬間拉滿的場(chǎng)景屬于典型的資源預(yù)算resource budgeting思路。pool_slots的實(shí)現(xiàn)與計(jì)費(fèi)邏輯可以在數(shù)據(jù)模型與調(diào)度器中交叉驗(yàn)證在 taskinstance.py 中pool_slots是任務(wù)實(shí)例上的整型列default1且不可為空任務(wù)實(shí)例創(chuàng)建時(shí)從任務(wù)定義上拷貝該值在 pool.py 的slots_stats中池的占用統(tǒng)計(jì)并不是“數(shù)任務(wù)個(gè)數(shù)”而是對(duì)處于執(zhí)行態(tài)的任務(wù)按池分組執(zhí)行func.sum(TaskInstance.pool_slots)——也就是說(shuō) heavy 任務(wù)在統(tǒng)計(jì)層面就被折算成了多個(gè)槽位相應(yīng)地occupied_slots()、running_slots()、queued_slots()等方法也都使用SUM(pool_slots)而非COUNT(*)。因此調(diào)度與展示兩個(gè)環(huán)節(jié)對(duì)“多槽位任務(wù)”的認(rèn)知是一致的一個(gè)pool_slots2的任務(wù)在統(tǒng)計(jì)、排隊(duì)、占坑全流程中都按 2 個(gè)單位計(jì)費(fèi)。調(diào)度器如何依據(jù)槽位放行任務(wù)源碼級(jí)原理理解了模型層的槽位計(jì)費(fèi)后再看調(diào)度器的具體決策邏輯可以完整還原“排隊(duì)—放行”的過(guò)程。核心實(shí)現(xiàn)在調(diào)度任務(wù)循環(huán) scheduler_job_runner.py 中調(diào)度器會(huì)先匯總當(dāng)前所有池的可用槽位并計(jì)算pool_slots_free如果沒(méi)有任何池還有空位則本輪的調(diào)度預(yù)算會(huì)被直接壓到 0對(duì)每個(gè)待調(diào)度的任務(wù)實(shí)例先讀取其所屬池的open_slots可用槽位若open_slots 0則本輪不放行記錄日志 “Not scheduling since there are 0 open slots in pool ...”若任務(wù)實(shí)例的pool_slots大于該池的總槽位pool_total說(shuō)明單任務(wù)所需的權(quán)重超過(guò)了池的容量上限任務(wù)不會(huì)被調(diào)度若任務(wù)實(shí)例的pool_slots大于當(dāng)前剩余open_slots同樣跳過(guò)等待槽位釋放每次放行一個(gè)任務(wù)調(diào)度器就執(zhí)行open_slots - task_instance.pool_slots更新該池在本輪迭代中的剩余容量供后續(xù)候選任務(wù)繼續(xù)判斷。其中“單任務(wù)pool_slots大于池總?cè)萘縿t不調(diào)度”的規(guī)則解釋了設(shè)計(jì)約束pool_slots應(yīng)該小于等于池的slots否則該任務(wù)永遠(yuǎn)無(wú)法獲得足夠槽位。另外調(diào)度器還會(huì)以pool.open_slots為指標(biāo)名把每個(gè)池的可用槽位上報(bào)到 metrics便于對(duì)池的擁堵程度做監(jiān)控告警。用 CLI 管理 Pools含 JSON 導(dǎo)入導(dǎo)出池不僅能在 UI 中管理也可以通過(guò)命令行腳本化維護(hù)便于把池的配置納入 IaC 流程。CLI 子命令在 cli_config.py 的POOLS_COMMANDS中定義實(shí)際實(shí)現(xiàn)位于 pool_command.py包括以下操作列出所有池airflow pools list可結(jié)合-o指定輸出格式如table、json、yaml。每條記錄會(huì)展示 pool 名、slots、description、include_deferred 與 team_name 字段。查看單個(gè)池airflow pools get pool_name池不存在時(shí)命令會(huì)以 “Pool ... does not exist” 退出。創(chuàng)建或更新池airflow pools set pool_name slots [description] [--include-deferred] [--team-name team_name]例如創(chuàng)建一個(gè)面向批處理管道的池airflow pools set ep_data_pipeline_db_msg_agg 10 DB message aggregation concurrency cap位置參數(shù)slots為整型決定并發(fā)上限--include-deferred控制是否把延遲任務(wù)計(jì)入占用--team-name用于把池歸屬到某個(gè)團(tuán)隊(duì)該選項(xiàng)需要 Airflow 開(kāi)啟multi_team模式在 pool.py 的create_or_update_pool中若未開(kāi)啟多團(tuán)隊(duì)模式而傳入team_name會(huì)直接拋出ValueError池已存在時(shí)執(zhí)行set會(huì)更新其槽位、描述與開(kāi)關(guān)即“不存在則創(chuàng)建、存在則更新”的冪等語(yǔ)義。刪除池airflow pools delete pool_name注意default_pool無(wú)法刪除刪除不存在的池會(huì)以 “Pool ... does not exist” 報(bào)錯(cuò)。從 JSON 文件導(dǎo)入池airflow pools import /path/to/pools.json導(dǎo)入文件支持的格式見(jiàn)ARG_POOL_IMPORT的幫助文本如下{ pool_1: {slots: 5, description: , include_deferred: true}, pool_2: {slots: 10, description: test, include_deferred: false, team_name: my_team} }將所有池導(dǎo)出到 JSON 文件airflow pools export /path/to/pools.json導(dǎo)出/導(dǎo)入組合非常適合在多個(gè)環(huán)境測(cè)試、預(yù)發(fā)、生產(chǎn)之間同步池配置。另外從源碼可以看到airflow pools系列命令在 pool_command.py 中標(biāo)注了deprecated_for_airflowctl(...)裝飾器提示其正逐步遷移到新的airflowctl管理入口如airflowctl pools list等在閱讀日志或遷移腳本時(shí)如遇到該提示屬于預(yù)期行為。實(shí)戰(zhàn)設(shè)計(jì)建議結(jié)合文檔與調(diào)度器行為給出幾條可直接落地的設(shè)計(jì)經(jīng)驗(yàn)為每個(gè)會(huì)被多 DAG 共享的下游系統(tǒng)建一個(gè)專屬池槽位數(shù)量以該系統(tǒng)實(shí)測(cè)可承受的峰值并發(fā)為準(zhǔn)而不是拍腦袋定大數(shù)重任務(wù)用pool_slots單獨(dú)計(jì)費(fèi)讓輕任務(wù)在重任務(wù)運(yùn)行期間仍有機(jī)會(huì)獲得剩余槽位或者反過(guò)來(lái)用重任務(wù)獨(dú)占容量來(lái)保護(hù)下游縮小default_pool從制度上促使每個(gè) DAG 作者顯式聲明資源邊界對(duì)池的queued任務(wù)堆積做監(jiān)控配合pool.open_slots指標(biāo)排隊(duì)時(shí)間異常增長(zhǎng)通常意味著容量不足或任務(wù)執(zhí)行時(shí)間超預(yù)期用airflow pools import/export把池配置版本化并在變更槽位時(shí)通過(guò)airflow pools set平滑調(diào)整避免重啟集群。延伸閱讀官方 Pools 文檔本文對(duì)應(yīng)原文延遲任務(wù)deferred tasks指南理解include_deferred開(kāi)關(guān)的作用對(duì)象Priority Weight 文檔排隊(duì)任務(wù)的放行順序規(guī)則Pool 數(shù)據(jù)模型與統(tǒng)計(jì)實(shí)現(xiàn)slot_pool表結(jié)構(gòu)、slots_stats/occupied_slots等槽位計(jì)算邏輯任務(wù)實(shí)例模型pool、pool_slots列及優(yōu)先級(jí)策略裝配調(diào)度器任務(wù)循環(huán)open_slots判定與pool_slots扣減的具體決策邏輯Pools CLI 命令定義 與 命令實(shí)現(xiàn)子命令、參數(shù)及 JSON 導(dǎo)入導(dǎo)出格式【免費(fèi)下載鏈接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows項(xiàng)目地址: https://gitcode.com/GitHub_Trending/ai/airflow創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考