斷點(diǎn)續(xù)傳:SSE、檢查點(diǎn)與冪等性實(shí)戰(zhàn)指南)
1. 項目概述當(dāng)AI生成任務(wù)“卡在”99%時做AI應(yīng)用開發(fā)尤其是涉及長文本生成、圖像繪制或者復(fù)雜推理鏈任務(wù)的兄弟估計都遇到過這種讓人血壓飆升的場景你精心設(shè)計了一個任務(wù)隊列用戶提交了一個需要耗時幾分鐘甚至更長的生成請求前端進(jìn)度條已經(jīng)走到了90%、95%甚至99%你都能想象用戶滿懷期待地盯著屏幕的樣子。然后啪連接斷了。可能是網(wǎng)絡(luò)波動可能是服務(wù)端某個實(shí)例重啟也可能是前端頁面被用戶不小心刷新了。結(jié)果就是用戶看到的是“生成失敗”或一片空白后臺那個已經(jīng)消耗了大量計算資源的任務(wù)就這么半途而廢了。更糟的是用戶重試系統(tǒng)又傻乎乎地從頭開始跑不僅浪費(fèi)資源用戶體驗也直接跌到谷底。這個項目要解決的就是這個“臨門一腳”的痛點(diǎn)。它的核心不是如何防止中斷網(wǎng)絡(luò)問題無法根除而是如何在中斷發(fā)生后讓任務(wù)能夠從斷點(diǎn)附近“接著跑”或者至少讓用戶能無縫地“接著看”把已完成的部分成果先拿到手把損失降到最低。這背后是三個關(guān)鍵技術(shù)思想的組合拳SSEServer-Sent Events用于實(shí)現(xiàn)實(shí)時的進(jìn)度流式推送讓用戶感知到“活著”檢查點(diǎn)Checkpoint用于在服務(wù)端持久化任務(wù)的中間狀態(tài)實(shí)現(xiàn)斷點(diǎn)續(xù)傳冪等Idempotency用于確保用戶重試操作的安全性避免重復(fù)執(zhí)行或狀態(tài)錯亂。接下來我就結(jié)合一個具體的AI長文本生成場景拆解這套方案的設(shè)計思路、實(shí)現(xiàn)細(xì)節(jié)和那些只有踩過坑才知道的注意事項。2. 核心架構(gòu)與設(shè)計思路拆解2.1 為什么是SSE而不是WebSocket或長輪詢首先我們需要一個機(jī)制讓服務(wù)器能主動、持續(xù)地向客戶端推送任務(wù)進(jìn)度。備選方案通常有WebSocket、長輪詢和SSE。WebSocket功能最強(qiáng)大全雙工通信。但對于我們這種典型的“服務(wù)器向客戶端單向推送進(jìn)度”的場景有點(diǎn)殺雞用牛刀。它引入的復(fù)雜度更高需要維護(hù)連接狀態(tài)、處理更復(fù)雜的協(xié)議對于只需要進(jìn)度推送的頁面來說不夠輕量。長輪詢Long Polling兼容性好但效率低。客戶端需要不斷地發(fā)起請求并掛起等待每次連接建立和斷開都有開銷在進(jìn)度頻繁更新的場景下顯得笨重。SSEServer-Sent Events它是基于HTTP的單向通信協(xié)議專門為服務(wù)器到客戶端的實(shí)時數(shù)據(jù)流設(shè)計。瀏覽器通過一個持久的HTTP連接監(jiān)聽服務(wù)器可以隨時通過這個連接推送數(shù)據(jù)片段。它的優(yōu)點(diǎn)非常契合我們的需求協(xié)議簡單基于HTTP/HTTPS無需額外端口或復(fù)雜協(xié)議握手防火墻友好。自動重連瀏覽器原生支持在連接斷開后自動嘗試重新連接并可以攜帶上一次接收到的事件ID這對于我們后續(xù)實(shí)現(xiàn)“斷點(diǎn)續(xù)傳”的感知非常有用。輕量級數(shù)據(jù)格式簡單data:、event:、id:等前綴開銷小。注意SSE有一個重要限制它是單向的服務(wù)器-客戶端。如果您的應(yīng)用在任務(wù)執(zhí)行過程中還需要頻繁接收客戶端的指令例如調(diào)整生成參數(shù)那么WebSocket可能是更好的選擇。但對于絕大多數(shù)“提交-等待-接收結(jié)果”的AI生成任務(wù)SSE的簡單和高效是首選。在我們的架構(gòu)里SSE通道主要負(fù)責(zé)推送兩種信息增量進(jìn)度如“已生成30%”和分塊的結(jié)果數(shù)據(jù)如生成文本的一段。這樣即使最終連接中斷用戶也已經(jīng)看到了大部分內(nèi)容體驗上的“斷裂感”會大大減輕。2.2 檢查點(diǎn)Checkpoint任務(wù)狀態(tài)的“存檔點(diǎn)”SSE解決了信息流的問題但任務(wù)本身在服務(wù)端的執(zhí)行狀態(tài)如何保存這就需要檢查點(diǎn)機(jī)制。想象一下玩單機(jī)游戲時的存檔功能檢查點(diǎn)就是任務(wù)的“存檔”。存什么這取決于任務(wù)類型。對于AI文本生成如使用大語言模型檢查點(diǎn)可能需要保存已生成的token序列、模型當(dāng)前的內(nèi)部狀態(tài)如Transformer的Key-Value緩存、生成參數(shù)temperature, top_p等、以及當(dāng)前進(jìn)度百分比。對于圖像生成則可能是擴(kuò)散模型去噪過程中的潛在特征、當(dāng)前步驟數(shù)等。存在哪需要是一個持久化、可快速讀寫的存儲。Redis是最常見的選擇因為它速度快支持復(fù)雜數(shù)據(jù)結(jié)構(gòu)如Hash。對于非常龐大的狀態(tài)如大型模型的完整緩存也可以考慮寫入對象存儲如S3/MinIO或數(shù)據(jù)庫但讀取恢復(fù)時會慢一些。一個折中方案是將核心元數(shù)據(jù)和輕量狀態(tài)存Redis大型二進(jìn)制狀態(tài)存對象存儲。何時存頻繁保存會影響性能太少則可能丟失過多進(jìn)度。常見的策略有周期性保存每生成N個token或每經(jīng)過M秒保存一次。里程碑式保存每完成一個邏輯段落或章節(jié)保存。結(jié)合SSE事件保存每次通過SSE推送進(jìn)度后異步觸發(fā)一次狀態(tài)保存。檢查點(diǎn)的存在使得工作進(jìn)程Celery worker或K8s Pod在意外退出后新啟動的進(jìn)程能夠從最新的檢查點(diǎn)加載狀態(tài)繼續(xù)執(zhí)行而不是從頭開始。2.3 冪等Idempotency防止“重復(fù)提交”的護(hù)城河當(dāng)連接斷開用戶下意識會點(diǎn)擊“重試”按鈕。如果沒有冪等控制這個操作可能會向消息隊列中插入一個新的任務(wù)導(dǎo)致同一個任務(wù)被執(zhí)行兩次浪費(fèi)資源甚至可能產(chǎn)生重復(fù)內(nèi)容或狀態(tài)沖突。冪等的核心是對于同一個操作無論執(zhí)行一次還是多次其產(chǎn)生的副作用應(yīng)該是一致的。對于我們的任務(wù)提交接口需要實(shí)現(xiàn)冪等。如何實(shí)現(xiàn)通常利用一個客戶端生成的唯一冪等鍵Idempotency Key。這個Key可以前端在創(chuàng)建任務(wù)時生成如UUID并隨任務(wù)請求一起發(fā)送。服務(wù)端處理流程服務(wù)端收到請求攜帶冪等鍵idempotency_key: req_abc123。以這個Key為索引查詢Redis或數(shù)據(jù)庫檢查是否已處理過相同Key的請求。如果未處理過執(zhí)行業(yè)務(wù)邏輯創(chuàng)建任務(wù)并將處理結(jié)果如任務(wù)ID、初始狀態(tài)與這個冪等鍵關(guān)聯(lián)存儲起來設(shè)置一個合理的過期時間如24小時。如果已處理過直接返回之前存儲的結(jié)果而不是創(chuàng)建新任務(wù)。對于“繼續(xù)任務(wù)”的請求同樣適用此邏輯確保同一個“繼續(xù)”請求不會重復(fù)觸發(fā)恢復(fù)流程。這樣即使用戶在斷線后瘋狂點(diǎn)擊重試也只有第一個請求會真正創(chuàng)建或恢復(fù)任務(wù)后續(xù)請求都只是返回既有的結(jié)果保證了系統(tǒng)的穩(wěn)定性和資源利用率。3. 系統(tǒng)核心組件與交互流程詳解3.1 任務(wù)生命周期與狀態(tài)機(jī)一個具備斷點(diǎn)續(xù)傳能力的AI生成任務(wù)其生命周期比普通任務(wù)更復(fù)雜。我們需要定義一個清晰的狀態(tài)機(jī)。以下是一個典型的狀態(tài)流轉(zhuǎn)設(shè)計PENDING等待中任務(wù)已創(chuàng)建進(jìn)入消息隊列等待Worker領(lǐng)取。PROCESSING處理中Worker開始執(zhí)行任務(wù)。此狀態(tài)需細(xì)分PROCESSING_STREAMING流式處理中核心生成邏輯正在運(yùn)行并通過SSE推送進(jìn)度和分片結(jié)果。同時定期寫入檢查點(diǎn)。PAUSED_BY_FAILURE因失敗暫停當(dāng)Worker進(jìn)程意外崩潰、任務(wù)執(zhí)行超時或遇到不可恢復(fù)錯誤時進(jìn)入此狀態(tài)。此時最新的檢查點(diǎn)已被保存。RECOVERING恢復(fù)中用戶發(fā)起“繼續(xù)”請求系統(tǒng)根據(jù)任務(wù)ID找到最新的檢查點(diǎn)啟動新的Worker實(shí)例加載狀態(tài)。COMPLETED已完成任務(wù)正常執(zhí)行完畢最終結(jié)果已持久化。FAILED已失敗任務(wù)遇到明確錯誤且無法恢復(fù)如輸入?yún)?shù)非法或重試次數(shù)耗盡。這個狀態(tài)機(jī)需要持久化在數(shù)據(jù)庫如tasks表中并且關(guān)鍵狀態(tài)變遷如從PROCESSING變?yōu)镻AUSED_BY_FAILURE需要記錄時間戳和原因便于排查問題。3.2 SSE服務(wù)端實(shí)現(xiàn)關(guān)鍵細(xì)節(jié)以Python的FastAPI框架為例實(shí)現(xiàn)SSE端點(diǎn)from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse import asyncio import json app FastAPI() # 一個全局字典用于管理任務(wù)與事件流隊列的映射生產(chǎn)環(huán)境應(yīng)用Redis等分布式結(jié)構(gòu) task_event_queues {} async def event_publisher(task_id: str): SSE事件發(fā)布器 queue asyncio.Queue() task_event_queues[task_id] queue try: while True: # 從隊列中獲取事件 event_data await queue.get() if event_data is None: # 收到終止信號 break # 格式化SSE數(shù)據(jù)data: 后跟JSON字符串結(jié)尾兩個換行 # 可以添加 event: progress 或 event: chunk 來區(qū)分事件類型 # 添加 id: sequence_id 支持客戶端斷線重連時指定Last-Event-ID yield fevent: {event_data[event]}\ndata: {json.dumps(event_data[data])}\nid: {event_data.get(id, )}\n\n finally: # 清理資源 task_event_queues.pop(task_id, None) app.get(/task/{task_id}/stream) async def stream_task_progress(task_id: str, request: Request): SSE流式端點(diǎn) async def generate(): async for message in event_publisher(task_id): yield message return StreamingResponse( generate(), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no # 對Nginx代理很重要禁用緩沖 } ) # 在任務(wù)Worker中推送進(jìn)度和分片 async def push_progress(task_id: str, progress: float, chunk: str None): queue task_event_queues.get(task_id) if queue: event { event: progress, data: {progress: progress, chunk: chunk}, id: fseq_{int(time.time()*1000)} # 簡單的事件ID } await queue.put(event)實(shí)操心得X-Accel-Buffering: no這個響應(yīng)頭對于通過Nginx等反向代理使用SSE至關(guān)重要。默認(rèn)情況下代理可能會緩沖數(shù)據(jù)導(dǎo)致客戶端無法實(shí)時收到消息。加上這個頭可以禁用代理緩沖。3.3 檢查點(diǎn)的存儲與恢復(fù)設(shè)計檢查點(diǎn)的存儲結(jié)構(gòu)需要精心設(shè)計。以下是一個基于Redis Hash的檢查點(diǎn)結(jié)構(gòu)示例import pickle import json import redis from typing import Any, Dict class CheckpointManager: def __init__(self, redis_client: redis.Redis): self.redis redis_client def save_checkpoint(self, task_id: str, checkpoint_data: Dict[str, Any], is_final: bool False): 保存檢查點(diǎn)。 :param is_final: 是否為最終結(jié)果如果是則檢查點(diǎn)可以標(biāo)記為完成態(tài)。 key fcheckpoint:{task_id} # 序列化數(shù)據(jù)。注意pickle可能有不兼容風(fēng)險對于簡單結(jié)構(gòu)推薦用JSON。 # 復(fù)雜對象如模型狀態(tài)可能需要特殊處理。 serialized_data pickle.dumps(checkpoint_data) pipeline self.redis.pipeline() pipeline.hset(key, data, serialized_data) pipeline.hset(key, updated_at, time.time()) pipeline.hset(key, progress, checkpoint_data.get(progress, 0.0)) if is_final: pipeline.hset(key, status, completed) else: pipeline.hset(key, status, in_progress) # 設(shè)置一個較長的過期時間比如7天確保用戶有足夠時間恢復(fù) pipeline.expire(key, 604800) pipeline.execute() def load_checkpoint(self, task_id: str) - Dict[str, Any]: key fcheckpoint:{task_id} data self.redis.hget(key, data) if not data: raise ValueError(fNo checkpoint found for task {task_id}) return pickle.loads(data) def get_checkpoint_info(self, task_id: str) - Dict: 獲取檢查點(diǎn)的元信息不加載完整數(shù)據(jù) key fcheckpoint:{task_id} return self.redis.hgetall(key)在Worker中的集成點(diǎn)# 在生成循環(huán)中 for i, chunk in enumerate(generator): # ... 處理生成邏輯 ... current_progress min(0.99, (i 1) / estimated_total_chunks) # 避免提前到100% # 推送SSE await push_progress(task_id, current_progress, chunk) # 每生成N個chunk或每過一段時間保存檢查點(diǎn) if i % checkpoint_interval 0: checkpoint_data { generated_tokens: all_tokens_so_far, model_state: generator.get_state(), # 假設(shè)模型有獲取狀態(tài)的方法 progress: current_progress, parameters: generation_params, } checkpoint_manager.save_checkpoint(task_id, checkpoint_data)4. 客戶端前端的協(xié)同與健壯性處理服務(wù)端再健壯也需要前端配合才能提供完整的用戶體驗。前端邏輯的核心是建立連接、處理流、管理重試、優(yōu)雅降級。4.1 建立SSE連接與處理事件class TaskStreamer { constructor(taskId) { this.taskId taskId; this.eventSource null; this.lastEventId null; // 用于斷線重連 this.retryCount 0; this.maxRetries 3; } connect() { if (this.eventSource) { this.eventSource.close(); } let url /api/task/${this.taskId}/stream; // 如果存在上次斷開時的最后事件ID可以帶上服務(wù)端可據(jù)此跳過已發(fā)送的事件 if (this.lastEventId) { url ?lastEventId${encodeURIComponent(this.lastEventId)}; } this.eventSource new EventSource(url); this.eventSource.addEventListener(progress, (event) { const data JSON.parse(event.data); this.lastEventId event.lastEventId || event.id; // 更新最后事件ID this.updateProgressUI(data.progress); if (data.chunk) { this.appendChunkToOutput(data.chunk); } }); this.eventSource.addEventListener(error, (event) { console.error(SSE connection error:, event); // 錯誤處理可能是網(wǎng)絡(luò)問題或服務(wù)端錯誤 this.handleStreamError(); }); this.eventSource.addEventListener(done, (event) { const data JSON.parse(event.data); console.log(Task completed:, data); this.eventSource.close(); this.showCompletionUI(data.finalResult); }); } handleStreamError() { this.retryCount; if (this.retryCount this.maxRetries) { console.log(Retrying connection... (${this.retryCount}/${this.maxRetries})); setTimeout(() this.connect(), 1000 * this.retryCount); // 指數(shù)退避 } else { // 重試失敗顯示錯誤并提供一個“手動恢復(fù)”按鈕 this.showRecoveryOption(); } } updateProgressUI(progress) { // 更新進(jìn)度條例如到95%時減緩動畫給用戶“即將完成”的心理預(yù)期 const progressBar document.getElementById(progress-bar); progressBar.style.width ${progress * 100}%; if (progress 0.95) { progressBar.classList.add(almost-done); // 添加一個緩慢脈沖的CSS動畫 } } appendChunkToOutput(chunk) { // 將流式生成的文本塊追加到顯示區(qū)域 // 注意對于Markdown或代碼可能需要特殊處理渲染 const outputEl document.getElementById(output); outputEl.innerHTML chunk; // 簡單追加生產(chǎn)環(huán)境考慮虛擬滾動 // 滾動到底部 outputEl.scrollTop outputEl.scrollHeight; } showRecoveryOption() { // 顯示一個按鈕“連接斷開點(diǎn)擊嘗試恢復(fù)任務(wù)” // 點(diǎn)擊后調(diào)用一個特定的“恢復(fù)任務(wù)”API而不是重新提交。 const recoveryBtn document.createElement(button); recoveryBtn.textContent 連接已斷開點(diǎn)擊恢復(fù)任務(wù); recoveryBtn.onclick () this.attemptRecovery(); document.getElementById(status-area).appendChild(recoveryBtn); } async attemptRecovery() { // 調(diào)用冪等的“恢復(fù)任務(wù)”接口 const resp await fetch(/api/task/${this.taskId}/recover, { method: POST, headers: { Idempotency-Key: recover_${this.taskId}_${Date.now()} // 生成冪等鍵 } }); if (resp.ok) { // 恢復(fù)成功重新建立SSE連接 this.retryCount 0; this.connect(); } else { alert(恢復(fù)任務(wù)失敗請稍后重試或聯(lián)系支持。); } } }4.2 頁面生命周期與連接管理前端必須妥善處理頁面刷新或關(guān)閉的情況。一個常見的策略是使用sessionStorage或localStorage來保存當(dāng)前任務(wù)的關(guān)鍵信息。// 頁面加載時 window.addEventListener(load, () { const savedTaskId sessionStorage.getItem(currentTaskId); const savedProgress sessionStorage.getItem(currentTaskProgress); if (savedTaskId) { // 提示用戶是否恢復(fù)上一個任務(wù) if (confirm(檢測到有未完成的任務(wù)是否恢復(fù))) { initializeTaskStreamer(savedTaskId, parseFloat(savedProgress)); } else { // 用戶選擇不恢復(fù)清理存儲 sessionStorage.removeItem(currentTaskId); sessionStorage.removeItem(currentTaskProgress); } } }); // 開始任務(wù)時 function startNewTask(params) { // ... 調(diào)用API創(chuàng)建任務(wù) ... const taskId response.task_id; sessionStorage.setItem(currentTaskId, taskId); sessionStorage.setItem(currentTaskProgress, 0); const streamer new TaskStreamer(taskId); streamer.connect(); } // 在streamer的updateProgressUI中更新存儲的進(jìn)度 updateProgressUI(progress) { // ... UI更新 ... sessionStorage.setItem(currentTaskProgress, progress.toString()); } // 任務(wù)完成或明確失敗時清理存儲 function cleanupTask() { sessionStorage.removeItem(currentTaskId); sessionStorage.removeItem(currentTaskProgress); }5. 服務(wù)端Worker的容錯與恢復(fù)實(shí)現(xiàn)Worker是執(zhí)行任務(wù)的核心它的健壯性直接決定了斷點(diǎn)續(xù)傳的可靠性。5.1 Worker的優(yōu)雅關(guān)閉與狀態(tài)保存無論是系統(tǒng)信號SIGTERM還是任務(wù)超時Worker都應(yīng)該有機(jī)會保存當(dāng)前狀態(tài)。import signal import sys from celery import Celery, Task from checkpoint_manager import CheckpointManager app Celery(ai_tasks) class ResilientTask(Task): 自定義Task基類增加優(yōu)雅關(guān)閉和檢查點(diǎn)支持 def __init__(self): super().__init__() self.checkpoint_manager CheckpointManager(redis_client) self._is_terminating False self._current_task_id None # 注冊信號處理器 signal.signal(signal.SIGTERM, self.handle_terminate) signal.signal(signal.SIGINT, self.handle_terminate) def handle_terminate(self, signum, frame): 收到終止信號時的處理 self._is_terminating True print(fReceived signal {signum}, attempting graceful shutdown...) # 注意這里不能做耗時操作只能設(shè)置標(biāo)志。真正的保存應(yīng)在任務(wù)循環(huán)中檢查此標(biāo)志。 def run(self, task_id, *args, **kwargs): self._current_task_id task_id # 1. 嘗試加載檢查點(diǎn)如果是恢復(fù)任務(wù) try: checkpoint self.checkpoint_manager.load_checkpoint(task_id) start_from_checkpoint True print(fResuming task {task_id} from checkpoint at progress {checkpoint.get(progress, 0)}) except ValueError: checkpoint None start_from_checkpoint False print(fStarting new task {task_id}) # 2. 任務(wù)主循環(huán) generator self.initialize_generator(checkpoint) for chunk_num, chunk in enumerate(generator): # 檢查終止標(biāo)志 if self._is_terminating: print(Termination signal received, saving checkpoint before exit...) self.save_current_state(task_id, generator, chunk_num) # 更新任務(wù)狀態(tài)為“因失敗暫停” update_task_status(task_id, PAUSED_BY_FAILURE, reasonworker_terminated) return {status: interrupted, message: Task paused due to worker shutdown.} # 正常業(yè)務(wù)邏輯處理chunk推送SSE... # ... # 定期保存檢查點(diǎn) if chunk_num % 10 0: self.save_current_state(task_id, generator, chunk_num) # 3. 任務(wù)完成 self.checkpoint_manager.save_checkpoint(task_id, {status: completed, final_result: ...}, is_finalTrue) update_task_status(task_id, COMPLETED) return {status: success, result: final_result} def save_current_state(self, task_id, generator, chunk_num): 保存當(dāng)前狀態(tài)到檢查點(diǎn) state { progress: chunk_num / estimated_total_chunks, generator_state: generator.get_state(), last_chunk: last_chunk, # ... 其他需要保存的狀態(tài) } self.checkpoint_manager.save_checkpoint(task_id, state)5.2 任務(wù)恢復(fù)API的實(shí)現(xiàn)當(dāng)用戶點(diǎn)擊“恢復(fù)”按鈕時前端調(diào)用的是一個專門的恢復(fù)接口而不是重新提交任務(wù)。from fastapi import APIRouter, HTTPException, Header from celery.result import AsyncResult import uuid router APIRouter() router.post(/task/{task_id}/recover) async def recover_task( task_id: str, idempotency_key: str Header(None, aliasIdempotency-Key) ): 冪等的任務(wù)恢復(fù)接口。 1. 檢查任務(wù)是否存在且狀態(tài)為 PAUSED_BY_FAILURE。 2. 檢查冪等鍵防止重復(fù)恢復(fù)。 3. 向隊列發(fā)送一個恢復(fù)任務(wù)而非普通任務(wù)。 # 1. 驗證任務(wù)狀態(tài) task_info get_task_from_db(task_id) if not task_info: raise HTTPException(status_code404, detailTask not found) if task_info.status ! PAUSED_BY_FAILURE: raise HTTPException(status_code400, detailfTask cannot be recovered from status {task_info.status}) # 2. 冪等性檢查 if idempotency_key: previous_response redis.get(fidempotency:{idempotency_key}) if previous_response: # 直接返回之前的響應(yīng) return json.loads(previous_response) else: # 如果客戶端沒提供可以生成一個但最好要求客戶端提供 idempotency_key fauto_{task_id}_{int(time.time())} # 3. 發(fā)送恢復(fù)任務(wù)到Celery。這里使用一個專門的RecoverTask。 # 與普通Task的區(qū)別在于RecoverTask會主動加載檢查點(diǎn)并從斷點(diǎn)開始執(zhí)行。 async_result app.send_task( tasks.recover_task, args[task_id], task_idfrecover_{task_id}_{uuid.uuid4().hex[:8]} # 生成新的Celery任務(wù)ID ) # 4. 更新原任務(wù)狀態(tài)為“恢復(fù)中” update_task_status(task_id, RECOVERING) # 5. 存儲冪等響應(yīng) response_data {recovery_task_id: async_result.id, status: recovery_initiated} redis.setex(fidempotency:{idempotency_key}, 86400, json.dumps(response_data)) # 24小時過期 return response_data6. 生產(chǎn)環(huán)境部署的注意事項與排查技巧6.1 網(wǎng)絡(luò)與代理配置SSE對網(wǎng)絡(luò)環(huán)境比較敏感尤其是在有反向代理Nginx, Apache, Cloudflare的情況下。Nginx配置必須禁用代理緩沖并設(shè)置合適的超時時間。location /api/task/ { proxy_pass http://backend_upstream; proxy_set_header Connection ; proxy_http_version 1.1; chunked_transfer_encoding off; proxy_buffering off; proxy_cache off; proxy_read_timeout 86400s; # 長連接超時時間根據(jù)任務(wù)時長調(diào)整 proxy_send_timeout 86400s; }Cloudflare默認(rèn)情況下Cloudflare會對響應(yīng)進(jìn)行緩沖可能中斷SSE。需要在Cloudflare的規(guī)則中為該SSE路徑設(shè)置“緩存級別”為“繞過”或者使用Cache-Control: no-cache, no-transform頭。6.2 檢查點(diǎn)數(shù)據(jù)的大小與性能AI模型的狀態(tài)可能非常大如LLM的KV緩存。全量保存和加載會帶來顯著的I/O開銷和內(nèi)存壓力。增量檢查點(diǎn)只保存自上次檢查點(diǎn)以來的變化部分。這對于某些模型結(jié)構(gòu)如RNN可能有效但對Transformer的完整KV緩存較難。壓縮對檢查點(diǎn)數(shù)據(jù)進(jìn)行壓縮如gzip后再存儲犧牲一些CPU時間換取網(wǎng)絡(luò)和存儲帶寬。分級存儲將核心元數(shù)據(jù)和小型狀態(tài)如進(jìn)度、參數(shù)存Redis將大型二進(jìn)制狀態(tài)如模型張量存對象存儲S3并在檢查點(diǎn)中保存引用指針。保存頻率權(quán)衡在“數(shù)據(jù)丟失風(fēng)險”和“性能開銷”之間取得平衡。對于耗時極長的任務(wù)如1小時以上可以每1-2分鐘保存一次對于短任務(wù)1分鐘內(nèi)可能只需要在開始和結(jié)束時保存或在中間保存一次。6.3 常見問題排查表問題現(xiàn)象可能原因排查步驟與解決方案前端收不到SSE事件1. 代理服務(wù)器緩沖。2. 服務(wù)端響應(yīng)頭不正確。3. 連接已斷開但前端未正確處理。1. 檢查Nginx等代理配置確認(rèn)已設(shè)置proxy_buffering off和X-Accel-Buffering: no。2. 瀏覽器開發(fā)者工具Network標(biāo)簽查看SSE連接確認(rèn)響應(yīng)頭Content-Type: text/event-stream連接狀態(tài)碼為200。3. 前端代碼添加EventSource的onerror和onopen監(jiān)聽器打印日志。任務(wù)恢復(fù)后進(jìn)度回退1. 檢查點(diǎn)保存頻率太低丟失了較多中間狀態(tài)。2. 恢復(fù)時加載了舊的檢查點(diǎn)。3. 模型狀態(tài)保存/恢復(fù)邏輯有誤。1. 增加檢查點(diǎn)保存頻率或在關(guān)鍵邏輯節(jié)點(diǎn)如段落結(jié)束強(qiáng)制保存。2. 檢查Redis中檢查點(diǎn)的updated_at時間戳確認(rèn)加載的是最新的。考慮使用版本號管理檢查點(diǎn)。3. 對比恢復(fù)前后模型的輸出是否具有連貫性。編寫單元測試驗證狀態(tài)序列化/反序列化的正確性。“恢復(fù)”按鈕點(diǎn)擊無效1. 冪等鍵沖突或處理邏輯有誤。2. 原任務(wù)狀態(tài)已不是PAUSED_BY_FAILURE。3. 恢復(fù)任務(wù)隊列堆積或Worker未啟動。1. 檢查服務(wù)端日志查看冪等鍵檢查邏輯。確保同一個恢復(fù)請求返回相同結(jié)果。2. 檢查數(shù)據(jù)庫中原任務(wù)的狀態(tài)字段。確保Worker崩潰時正確更新了狀態(tài)。3. 檢查Celery Worker日志確認(rèn)recover_task任務(wù)被正確消費(fèi)和執(zhí)行。內(nèi)存占用過高1. 同時保存過多任務(wù)的檢查點(diǎn)數(shù)據(jù)在內(nèi)存中。2. 模型狀態(tài)過大未及時清理。1. 為檢查點(diǎn)Redis實(shí)例設(shè)置內(nèi)存上限和淘汰策略如allkeys-lru。2. 實(shí)現(xiàn)檢查點(diǎn)的自動過期清理如任務(wù)完成后24小時刪除。對于大型狀態(tài)考慮使用外部存儲。瀏覽器標(biāo)簽頁關(guān)閉后無法恢復(fù)前端未在頁面卸載前保存足夠信息到持久化存儲。使用window.addEventListener(beforeunload, ...)事件在頁面關(guān)閉前將任務(wù)ID和進(jìn)度同步到localStorage比sessionStorage生命周期更長。下次打開同源頁面時檢查并提示恢復(fù)。6.4 監(jiān)控與告警一個健壯的系統(tǒng)離不開監(jiān)控。SSE連接數(shù)監(jiān)控活躍的SSE連接數(shù)異常升高可能預(yù)示連接泄漏或用戶激增。檢查點(diǎn)保存失敗率監(jiān)控保存檢查點(diǎn)到Redis/S3的失敗次數(shù)失敗可能意味著存儲服務(wù)異常或數(shù)據(jù)過大。任務(wù)中斷率計算狀態(tài)變?yōu)镻AUSED_BY_FAILURE的任務(wù)占總?cè)蝿?wù)的比例。如果比例突然升高需要檢查Worker健康狀況或底層基礎(chǔ)設(shè)施如GPU節(jié)點(diǎn)。任務(wù)恢復(fù)成功率監(jiān)控恢復(fù)請求的成功率過低可能意味著檢查點(diǎn)機(jī)制或恢復(fù)邏輯存在缺陷。端到端進(jìn)度感知可以在前端SSE事件中埋點(diǎn)上報“最后收到進(jìn)度的時間”。如果大量任務(wù)進(jìn)度長時間不更新但連接未斷可能意味著后端生成邏輯卡住。這套“SSE 檢查點(diǎn) 冪等”的組合方案本質(zhì)上是在分布式、不可靠的網(wǎng)絡(luò)環(huán)境中為長耗時任務(wù)構(gòu)建一個“韌性層”。它承認(rèn)失敗會發(fā)生但致力于讓失敗的影響變得可感知、可恢復(fù)、對用戶透明。在實(shí)際部署中根據(jù)具體的AI模型、業(yè)務(wù)負(fù)載和基礎(chǔ)設(shè)施每個環(huán)節(jié)都需要進(jìn)行細(xì)致的調(diào)優(yōu)和測試。例如對于超長文本生成可能需要結(jié)合“分頁”或“分段”生成策略將一個大任務(wù)拆分成多個可獨(dú)立保存檢查點(diǎn)的子任務(wù)進(jìn)一步降低單次故障的影響范圍。