實(shí)踐與源碼閱讀)
文章目錄從“能運(yùn)行”到“能生產(chǎn)”還差什么使用 Celery Beat 執(zhí)行周期任務(wù)為什么同一套計(jì)劃通常只能運(yùn)行一個(gè) Beat為什么需要多個(gè)隊(duì)列Worker 并發(fā)池怎么選擇PreforkEventlet / GeventSoloThreads并發(fā)數(shù)不是越大越好理解 Prefetch管理 Worker 子進(jìn)程生命周期優(yōu)雅停止與滾動(dòng)發(fā)布使用命令行觀察 Celery使用 Flower 進(jìn)行 Web 監(jiān)控真正應(yīng)該監(jiān)控哪些指標(biāo)隊(duì)列指標(biāo)任務(wù)指標(biāo)Worker 指標(biāo)依賴指標(biāo)生產(chǎn)環(huán)境安全配置Broker 和 Backend序列化敏感參數(shù)生產(chǎn)配置示例常見生產(chǎn)故障隊(duì)列持續(xù)積壓任務(wù)重復(fù)執(zhí)行Worker 內(nèi)存不斷增長(zhǎng)任務(wù)永遠(yuǎn)卡住如何測(cè)試生產(chǎn)行為單元測(cè)試集成測(cè)試故障演練從 GitHub 倉(cāng)庫(kù)理解 Celery應(yīng)用與配置celery/app/base.py任務(wù)對(duì)象celery/app/task.py結(jié)果抽象celery/result.py工作流celery/canvas.pyWorkercelery/worker/Result Backendcelery/backends/消息傳輸為什么經(jīng)常出現(xiàn) Kombu推薦源碼閱讀順序生產(chǎn)上線檢查清單架構(gòu)可靠性安全運(yùn)維系列總結(jié)參考資料從“能運(yùn)行”到“能生產(chǎn)”還差什么開發(fā)環(huán)境里一條命令啟動(dòng) Redis一條命令啟動(dòng) Worker任務(wù)成功返回似乎已經(jīng)完成。生產(chǎn)環(huán)境還必須回答周期任務(wù)如何避免重復(fù)調(diào)度視頻轉(zhuǎn)碼為什么不能和通知任務(wù)共用同一隊(duì)列Worker 并發(fā)數(shù)應(yīng)該設(shè)置多少隊(duì)列積壓時(shí)如何發(fā)現(xiàn)Worker 發(fā)布重啟時(shí)在途任務(wù)怎么辦任務(wù)參數(shù)是否泄漏敏感數(shù)據(jù)Broker 或 Backend 故障后如何恢復(fù)如何定位任務(wù)慢在排隊(duì)還是執(zhí)行Celery 是分布式系統(tǒng)的一部分不是一個(gè)裝飾器庫(kù)。生產(chǎn)化的重點(diǎn)是資源隔離、可靠性、安全和可觀測(cè)性。使用 Celery Beat 執(zhí)行周期任務(wù)Celery Beat 是調(diào)度器。它按計(jì)劃創(chuàng)建任務(wù)消息并發(fā)送到 Broker真正執(zhí)行任務(wù)的仍然是 Worker。配置固定間隔app.conf.beat_schedule{build-health-report-every-5-minutes:{task:tasks.build_health_report,schedule:300.0,},}配置 Crontabfromcelery.schedulesimportcrontab app.conf.beat_schedule{cleanup-every-night:{task:tasks.cleanup_expired_data,schedule:crontab(hour2,minute0),options:{queue:maintenance,expires:3600,},},}啟動(dòng) Workercelery-Acelery_app worker--loglevelINFO另開進(jìn)程啟動(dòng) Beatcelery-Acelery_app beat--loglevelINFO開發(fā)環(huán)境可以使用worker -B合并啟動(dòng)但生產(chǎn)環(huán)境更適合分開管理和擴(kuò)縮容。為什么同一套計(jì)劃通常只能運(yùn)行一個(gè) Beat如果兩個(gè) Beat 同時(shí)加載相同時(shí)間表它們可能在同一時(shí)刻各發(fā)送一次任務(wù)于是周期任務(wù)重復(fù)執(zhí)行。解決思路確保只有一個(gè) Beat 實(shí)例使用支持鎖或高可用選主的調(diào)度方案即使調(diào)度層防重任務(wù)本身仍保持冪等對(duì)必須單實(shí)例執(zhí)行的任務(wù)增加分布式鎖或業(yè)務(wù)狀態(tài)約束。還要考慮任務(wù)重疊每五分鐘調(diào)度一次但任務(wù)需要十分鐘下一次觸發(fā)時(shí)上一次還沒結(jié)束??梢允褂没跇I(yè)務(wù)鍵的鎖lock_key periodic:daily-settlement:2026-08-19鎖需要設(shè)置合理過期時(shí)間并處理 Worker 崩潰、鎖續(xù)期和誤釋放。很多場(chǎng)景下數(shù)據(jù)庫(kù)唯一約束比單純 Redis 鎖更容易形成可審計(jì)結(jié)果。為什么需要多個(gè)隊(duì)列假設(shè)同一隊(duì)列里同時(shí)存在50 毫秒的通知任務(wù)5 秒的第三方 API 調(diào)用30 分鐘的視頻轉(zhuǎn)碼高內(nèi)存的報(bào)表任務(wù)。長(zhǎng)任務(wù)占滿 Worker 后用戶通知會(huì)長(zhǎng)時(shí)間排隊(duì)高內(nèi)存任務(wù)還可能導(dǎo)致執(zhí)行其他任務(wù)的子進(jìn)程一起受到資源壓力。按工作負(fù)載分隊(duì)列app.conf.task_routes{tasks.send_email:{queue:io_fast},tasks.call_partner_api:{queue:io_external},tasks.transcode_video:{queue:cpu_heavy},tasks.build_report:{queue:memory_heavy},}分別啟動(dòng) Workercelery-Acelery_app worker\-Qio_fast\--concurrency20\--loglevelINFO celery-Acelery_app worker\-Qcpu_heavy\--concurrency4\--loglevelINFO資源隔離的收益長(zhǎng)任務(wù)不再阻塞短任務(wù)不同隊(duì)列可以獨(dú)立擴(kuò)縮容并發(fā)模型和資源限制可以分別配置單一業(yè)務(wù)故障不容易拖垮所有后臺(tái)任務(wù)隊(duì)列積壓更容易定位到具體工作負(fù)載。Worker 并發(fā)池怎么選擇Prefork默認(rèn)且最常用的多進(jìn)程模型。適合普通 Python 任務(wù)和 CPU 密集型工作進(jìn)程隔離也更明確。代價(jià)是每個(gè)子進(jìn)程都有內(nèi)存開銷創(chuàng)建大量進(jìn)程會(huì)增加數(shù)據(jù)庫(kù)連接和系統(tǒng)資源消耗。Eventlet / Gevent適合大量 I/O 等待且依賴庫(kù)能夠配合協(xié)作式并發(fā)的任務(wù)。需要 monkey patch并非所有庫(kù)都兼容。不要僅因?yàn)椤安l(fā)數(shù)可以設(shè)置很大”就使用。下游服務(wù)、數(shù)據(jù)庫(kù)連接池和限流策略仍然決定真實(shí)容量。Solo在主進(jìn)程單線程執(zhí)行適合調(diào)試或特殊環(huán)境沒有并行能力。Threads線程池可用于部分 I/O 場(chǎng)景但受 Python 庫(kù)線程安全性和 GIL 等因素影響需要基準(zhǔn)測(cè)試。并發(fā)數(shù)不是越大越好并發(fā)數(shù)受到多個(gè)瓶頸約束Worker 并發(fā) ≤ CPU / 內(nèi)存能力 ≤ 數(shù)據(jù)庫(kù)連接池 ≤ Redis / RabbitMQ 容量 ≤ 第三方 API 限流 ≤ 下游服務(wù)可承受并發(fā)如果數(shù)據(jù)庫(kù)只允許二十個(gè)連接卻啟動(dòng)一百個(gè)同時(shí)訪問數(shù)據(jù)庫(kù)的任務(wù)結(jié)果可能是更多超時(shí)和重試而不是更高吞吐。正確方法測(cè)量單任務(wù) CPU、內(nèi)存、I/O 和執(zhí)行時(shí)間確定下游容量從保守并發(fā)開始?jí)簻y(cè)觀察吞吐、錯(cuò)誤率和尾延遲按隊(duì)列分別調(diào)整。理解 PrefetchWorker 可以提前從 Broker 預(yù)取任務(wù)。預(yù)取能提高吞吐但也可能造成任務(wù)分配不均某個(gè) Worker 預(yù)取了大量長(zhǎng)任務(wù)其他 Worker 卻沒有工作。常見配置app.conf.worker_prefetch_multiplier1較低預(yù)取通常更適合長(zhǎng)任務(wù)和公平分配短小、穩(wěn)定的任務(wù)可能從更高預(yù)取獲得吞吐收益。worker_prefetch_multiplier1不是萬能最佳值。應(yīng)按隊(duì)列特征壓測(cè)。管理 Worker 子進(jìn)程生命周期第三方庫(kù)可能緩慢泄漏內(nèi)存。Celery 可以在子進(jìn)程處理一定任務(wù)數(shù)或達(dá)到內(nèi)存閾值后替換它app.conf.update(worker_max_tasks_per_child1000,worker_max_memory_per_child512_000,)含義子進(jìn)程最多執(zhí)行一千個(gè)任務(wù)后重啟子進(jìn)程內(nèi)存超過約 512 MB 后被替換。這些配置只能緩解問題不能代替定位內(nèi)存泄漏。頻繁重啟也會(huì)帶來初始化開銷。優(yōu)雅停止與滾動(dòng)發(fā)布Worker 收到TERM時(shí)會(huì)進(jìn)行溫和關(guān)閉停止接收新工作并等待當(dāng)前任務(wù)完成。QUIT更接近冷關(guān)閉SIGKILL則不給進(jìn)程清理機(jī)會(huì)。生產(chǎn)發(fā)布應(yīng)使用 systemd、Supervisor、Kubernetes 等管理進(jìn)程配置足夠長(zhǎng)的終止寬限時(shí)間停止前讓負(fù)載均衡或隊(duì)列逐步摘除 Worker觀察在途任務(wù)和隊(duì)列積壓對(duì)長(zhǎng)任務(wù)使用冪等和晚確認(rèn)時(shí)驗(yàn)證重投行為避免所有 Worker 同時(shí)退出。如果 Kubernetes 的terminationGracePeriodSeconds小于任務(wù)正常耗時(shí)所謂優(yōu)雅停止實(shí)際仍會(huì)變成強(qiáng)制終止。使用命令行觀察 Celerycelery-Acelery_app status celery-Acelery_app inspect registered celery-Acelery_app inspect active celery-Acelery_app inspect reserved celery-Acelery_app inspect scheduled celery-Acelery_app inspect stats含義registeredWorker 注冊(cè)了哪些任務(wù)active正在執(zhí)行reserved已被 Worker 預(yù)取但尚未執(zhí)行scheduledWorker 內(nèi)部等待 ETA 的任務(wù)stats進(jìn)程池、Broker 和運(yùn)行統(tǒng)計(jì)。這些命令依賴 Broker 對(duì)遠(yuǎn)程控制的支持。SQS 等 Broker 的能力與 RabbitMQ、Redis 不同。使用 Flower 進(jìn)行 Web 監(jiān)控Celery 官方監(jiān)控指南推薦 Flower 作為實(shí)時(shí) Web 監(jiān)控工具。安裝并啟動(dòng)pipinstallflower celery-Acelery_app flower--port5555訪問http://localhost:5555Flower 可以顯示W(wǎng)orker 在線狀態(tài)任務(wù)歷史、參數(shù)、狀態(tài)和運(yùn)行時(shí)間活躍、保留、計(jì)劃和撤銷任務(wù)Worker 池大小和隊(duì)列部分遠(yuǎn)程控制功能Prometheus 指標(biāo)集成。不要把 Flower 無認(rèn)證地暴露到公網(wǎng)。它可能顯示敏感任務(wù)參數(shù)并具備管理 Worker 和撤銷任務(wù)的能力。真正應(yīng)該監(jiān)控哪些指標(biāo)隊(duì)列指標(biāo)隊(duì)列長(zhǎng)度最老消息等待時(shí)間入隊(duì)和出隊(duì)速率未確認(rèn)消息數(shù)量各隊(duì)列消費(fèi)者數(shù)量。只看隊(duì)列長(zhǎng)度不夠。如果任務(wù)進(jìn)入和處理速度都很高隊(duì)列長(zhǎng)度可能穩(wěn)定最老消息年齡更能說明用戶等待多久。任務(wù)指標(biāo)成功率、失敗率和重試率P50、P95、P99 排隊(duì)時(shí)間P50、P95、P99 執(zhí)行時(shí)間超時(shí)和撤銷數(shù)量按任務(wù)類型統(tǒng)計(jì)的異常最終失敗和人工補(bǔ)償數(shù)量。Worker 指標(biāo)在線 Worker 和心跳CPU、內(nèi)存、負(fù)載和文件句柄子進(jìn)程異常退出和重啟當(dāng)前并發(fā)使用率Broker 重連次數(shù)。依賴指標(biāo)Broker 連接、內(nèi)存和磁盤Result Backend 延遲與容量數(shù)據(jù)庫(kù)連接池第三方 API 延遲、限流和錯(cuò)誤率對(duì)象存儲(chǔ)吞吐。Worker 在線不代表系統(tǒng)健康。Worker 全部在線但隊(duì)列最老消息已經(jīng)等待一小時(shí)業(yè)務(wù)仍然不可用。生產(chǎn)環(huán)境安全配置Broker 和 Backend使用獨(dú)立賬號(hào)和最小權(quán)限限制網(wǎng)絡(luò)訪問范圍開啟 TLS定期輪換憑據(jù)不與不可信應(yīng)用共享同一隊(duì)列或 Redis 數(shù)據(jù)庫(kù)按重要性設(shè)計(jì)持久化、備份和高可用。序列化默認(rèn)優(yōu)先 JSONapp.conf.update(task_serializerjson,result_serializerjson,accept_content[json],)pickle能表達(dá)更多 Python 類型但反序列化不可信 Pickle 數(shù)據(jù)可能執(zhí)行任意代碼。除非整個(gè)生產(chǎn)者、Broker 和 Worker 的信任邊界都經(jīng)過嚴(yán)格控制否則不要啟用。敏感參數(shù)任務(wù)參數(shù)可能出現(xiàn)在Broker 消息Worker 日志Flower監(jiān)控事件Result Backend異常追蹤系統(tǒng)。不要直接傳密碼、完整銀行卡號(hào)和訪問令牌。傳安全存儲(chǔ)中的引用由 Worker 在執(zhí)行時(shí)按權(quán)限讀取。使用argsrepr或kwargsrepr可以隱藏日志展示但不會(huì)加密 Broker 中的原始消息。生產(chǎn)配置示例importosfromceleryimportCelery appCelery(production_app,brokeros.environ[CELERY_BROKER_URL],backendos.environ[CELERY_RESULT_BACKEND],)app.conf.update(task_serializerjson,accept_content[json],result_serializerjson,enable_utcTrue,timezoneAsia/Singapore,result_expires3600,broker_connection_retry_on_startupTrue,worker_prefetch_multiplier1,worker_max_tasks_per_child1000,task_soft_time_limit300,task_time_limit330,task_routes{myapp.tasks.send_email:{queue:io_fast},myapp.tasks.build_report:{queue:reports},},)這只是起點(diǎn)不是所有系統(tǒng)通用的最佳配置。特別是預(yù)取、時(shí)間限制、進(jìn)程回收和結(jié)果過期時(shí)間必須根據(jù)任務(wù)特征測(cè)試。常見生產(chǎn)故障隊(duì)列持續(xù)積壓分析順序入隊(duì)速率是否突然增加Worker 數(shù)量或并發(fā)是否下降單任務(wù)執(zhí)行時(shí)間是否變長(zhǎng)下游數(shù)據(jù)庫(kù)或 API 是否變慢重試是否造成消息放大某類長(zhǎng)任務(wù)是否占滿共享隊(duì)列預(yù)取是否造成分配不均。不要第一反應(yīng)只擴(kuò)容 Worker。下游已經(jīng)飽和時(shí)擴(kuò)容會(huì)讓故障更嚴(yán)重。任務(wù)重復(fù)執(zhí)行檢查是否使用acks_lateWorker 是否在執(zhí)行中失聯(lián)Redis visibility timeout 是否短于任務(wù)耗時(shí)是否啟動(dòng)了多個(gè) Beat生產(chǎn)者是否因 HTTP 重試重復(fù)發(fā)送任務(wù)是否缺少業(yè)務(wù)冪等鍵。Worker 內(nèi)存不斷增長(zhǎng)檢查任務(wù)是否加載超大數(shù)據(jù)庫(kù)或全局緩存是否泄漏返回值是否過大Prefork 子進(jìn)程是否長(zhǎng)期不回收是否可以流式處理或分塊worker_max_tasks_per_child能否臨時(shí)緩解。任務(wù)永遠(yuǎn)卡住優(yōu)先尋找沒有超時(shí)的網(wǎng)絡(luò)請(qǐng)求數(shù)據(jù)庫(kù)鎖等待子進(jìn)程或外部命令未設(shè)置超時(shí)無限循環(huán)在任務(wù)里調(diào)用其他任務(wù)的.get()。如何測(cè)試生產(chǎn)行為單元測(cè)試把業(yè)務(wù)邏輯和 Celery 外殼分開defcalculate_invoice(order_id:int)-dict:...app.task(autoretry_for(TemporaryError,),retry_backoffTrue)defcalculate_invoice_task(order_id:int)-dict:returncalculate_invoice(order_id)普通函數(shù)可以快速、穩(wěn)定地單元測(cè)試。集成測(cè)試啟動(dòng)真實(shí)測(cè)試 Broker 和 Worker驗(yàn)證任務(wù)注冊(cè)JSON 序列化路由重試結(jié)果存儲(chǔ)Chain、Group 和 ChordWorker 退出后的行為。task_always_eagerTrue在當(dāng)前進(jìn)程同步執(zhí)行不能覆蓋 Broker、Worker、并發(fā)和消息確認(rèn)因此不能代替集成測(cè)試。故障演練主動(dòng)測(cè)試Broker 短暫斷開Worker 執(zhí)行中被終止下游服務(wù)超時(shí)和限流Result Backend 不可用隊(duì)列突然積壓Beat 重啟同一任務(wù)重復(fù)投遞。只有在故障中驗(yàn)證過的恢復(fù)方案才接近可信。從 GitHub 倉(cāng)庫(kù)理解 CeleryCelery 的官方主倉(cāng)庫(kù)是celery/celery。閱讀源碼時(shí)不建議從 Worker 啟動(dòng)流程一路硬追到底而應(yīng)圍繞已經(jīng)理解的概念分層閱讀。應(yīng)用與配置celery/app/base.pycelery/app/base.py包含核心Celery應(yīng)用對(duì)象。重點(diǎn)搜索class Celerysend_task配置加載Backend 與連接創(chuàng)建任務(wù)注冊(cè)和自動(dòng)發(fā)現(xiàn)。它回答“Celery 應(yīng)用怎樣把配置、任務(wù)和通信能力組織在一起”。任務(wù)對(duì)象celery/app/task.pycelery/app/task.py是理解使用層行為的關(guān)鍵。重點(diǎn)搜索class Taskdelayapply_asyncretry__call__生命周期鉤子。這里可以看清task(1,2)task.delay(1,2)為什么走的是完全不同的路徑。結(jié)果抽象celery/result.pycelery/result.py包含AsyncResult、GroupResult等。一個(gè)重要認(rèn)知是AsyncResult自身不是存放最終結(jié)果的容器它是根據(jù)任務(wù) ID 查詢 Result Backend 的抽象。工作流celery/canvas.pycelery/canvas.py實(shí)現(xiàn) Signature、Chain、Group、Chord 等 Canvas 原語(yǔ)。先讀官方 Canvas 文檔再結(jié)合源碼查找同名類和方法會(huì)比直接讀整份文件更高效。Workercelery/worker/celery/worker涵蓋 Worker、Consumer、并發(fā)池交互和啟動(dòng)組件。建議帶著問題閱讀Worker 如何連接 BrokerConsumer 如何接收消息收到消息后怎樣轉(zhuǎn)成任務(wù)請(qǐng)求請(qǐng)求如何交給進(jìn)程池成功、失敗、重試和確認(rèn)分別在哪里發(fā)生Result Backendcelery/backends/celery/backends包含 Redis、數(shù)據(jù)庫(kù)等結(jié)果后端實(shí)現(xiàn)和共同抽象。當(dāng)遇到 Chord、結(jié)果過期或 Backend 連接問題時(shí)這一目錄很有價(jià)值。消息傳輸為什么經(jīng)常出現(xiàn) KombuCelery 通過 Kombu抽象 RabbitMQ、Redis、SQS 等消息傳輸。因此報(bào)錯(cuò)棧中常出現(xiàn)kombu.connection kombu.transport.redis kombu.messagingCelery 負(fù)責(zé)任務(wù)語(yǔ)義與執(zhí)行Kombu 負(fù)責(zé)更底層的消息連接、Producer、Consumer 和 Transport 抽象。推薦源碼閱讀順序閱讀主倉(cāng)庫(kù) README明確項(xiàng)目邊界和支持環(huán)境在app/task.py中跟蹤delay → apply_async在app/base.py中閱讀send_task在result.py中理解AsyncResult在canvas.py中對(duì)應(yīng) Signature、Chain、Group、Chord在worker/consumer/中跟蹤消息消費(fèi)閱讀具體 Backend最后進(jìn)入 Kombu 查看所用 Broker 的 Transport。閱讀方法使用rg def apply_async celery/搜索入口用調(diào)試器或日志驗(yàn)證調(diào)用鏈固定 Celery 版本不要拿 main 分支源碼解釋舊版本生產(chǎn)行為先回答一個(gè)具體問題再擴(kuò)展閱讀范圍結(jié)合單元測(cè)試?yán)斫膺吔缜闆r。生產(chǎn)上線檢查清單架構(gòu)Broker、Backend 的角色和容量明確長(zhǎng)短任務(wù)、CPU 與 I/O 任務(wù)已經(jīng)分隊(duì)列各隊(duì)列有獨(dú)立擴(kuò)縮容策略Beat 單實(shí)例或具備可靠選主關(guān)鍵任務(wù)具有業(yè)務(wù)冪等方案??煽啃灾恢卦嚳苫謴?fù)異常配置指數(shù)退避、jitter 和最大次數(shù)所有網(wǎng)絡(luò) I/O 有超時(shí)明確提前確認(rèn)或晚確認(rèn)數(shù)據(jù)庫(kù)事務(wù)與消息發(fā)送不存在明顯競(jìng)態(tài)永久失敗可告警、查詢和補(bǔ)償。安全Broker 和 Backend 使用最小權(quán)限憑據(jù)由密鑰系統(tǒng)或環(huán)境變量提供網(wǎng)絡(luò)訪問受到限制并啟用 TLS只接受可信序列化格式任務(wù)參數(shù)不包含敏感明文Flower 有認(rèn)證且不直接暴露公網(wǎng)。運(yùn)維Worker 支持優(yōu)雅停止發(fā)布寬限時(shí)間覆蓋合理任務(wù)時(shí)長(zhǎng)監(jiān)控隊(duì)列長(zhǎng)度和最老消息年齡監(jiān)控成功率、重試率和尾延遲監(jiān)控 Worker 心跳、CPU 和內(nèi)存Broker、Backend 和下游依賴都有告警已完成 Worker、Broker 和下游故障演練。系列總結(jié)學(xué)會(huì) Celery 可以分為五個(gè)層次理解角色生產(chǎn)者、Broker、Worker、Backend 和 Beat跑通鏈路定義任務(wù)、啟動(dòng) Worker、發(fā)送消息、讀取結(jié)果保證正確重試、超時(shí)、確認(rèn)、重復(fù)執(zhí)行和冪等表達(dá)流程用 Canvas 組合順序、并行和匯總?cè)蝿?wù)長(zhǎng)期運(yùn)行隊(duì)列隔離、容量、監(jiān)控、安全、發(fā)布和故障恢復(fù)。最值得記住的一句話是Celery 負(fù)責(zé)可靠地分發(fā)和執(zhí)行任務(wù)但業(yè)務(wù)是否正確最終仍取決于冪等性、事務(wù)邊界、資源治理和可觀測(cè)性的設(shè)計(jì)。參考資料Celery 5.6 官方文檔Periodic TasksRouting TasksWorkers GuideMonitoring and Management GuideSecurityOptimizingcelery/celery GitHub 主倉(cāng)庫(kù)celery/kombu GitHub 倉(cāng)庫(kù)